MySQL Binlog消息接入实战教程
前言
本文以「业务表变更 → Binlog → 消息队列 → 应用消费」的增量同步链路为例,讲清楚 Binlog 消息 从哪来、长什么样、怎么发出来、怎么被消费。文中所有字段、表名均为演示用途,给出的是一套可直接落地的标准做法。
一、Binlog 消息从哪来
1、什么是 Binlog
MySQL 的 Binary Log(二进制日志) 记录了所有对数据产生变更的操作(INSERT / UPDATE / DELETE / DDL),最初用于主从复制和数据恢复。
利用这一特性,我们可以让一个「订阅者」伪装成 MySQL 的从库,实时拉取主库的 Binlog,从而感知业务表的每一次数据变更——这就是 CDC(Change Data Capture,变更数据捕获)。
2、完整链路
业务系统 (写库)
│ INSERT / UPDATE / DELETE
▼
MySQL 主库 ──写入──► Binlog
│
▼
采集组件 (Canal / Debezium / DTS 等) ← 伪装成 MySQL 从库,订阅 Binlog
│ 解析 Binlog → 转成 JSON
▼
消息队列 (Kafka / RocketMQ 等) ← 按表/库路由到指定 Topic
│
▼
本应用消费 (消息监听器) ← 本文重点
│ 反序列化 → 分发 → 业务处理
▼
下游 (缓存 / ES / 增量表 / 派生数据)
关键点:应用不直接连接 MySQL,而是消费采集组件投递到消息队列的变更消息。这样做的好处:
- 解耦:业务写库无感知,不影响主流程。
- 削峰:消息队列缓冲高峰流量。
- 可回溯:消息可重放,支持重试。
二、Binlog 消息长什么样(报文结构)
采集组件把一行数据的变更序列化为一条 JSON 消息,投递到消息队列。一个通用且清晰的报文结构如下:
{
"version": "1.0",
"database": "demo_db",
"table": "order_status",
"type": "update",
"ts": 1725850000000,
"pks": "id",
"data": {
"id": "1001",
"biz_code": "A-2026",
"biz_type": "NORMAL",
"status": "PAID",
"modified_time": "2026-09-09 10:00:00"
},
"old": {
"status": "UNPAID"
}
}
| 字段 | 含义 |
|---|---|
version |
报文版本号 |
database |
库名 |
table |
表名(分发的关键依据) |
type |
操作类型:insert / update / delete / ddl |
ts |
变更时间戳(毫秒) |
pks |
主键字段名 |
data |
变更后的整行数据(after) |
old |
变更前的字段数据(before,通常只含被改动字段) |
格式差异怎么办:不同采集组件的字段命名可能有差异(如有的用
tab表示表名、cur表示变更后)。消费端不应把某一种格式写死,而是在入口做一层「格式归一化」——用策略模式把外部报文统一转成内部标准模型,后续逻辑只面向标准模型编程:
/** 标准领域事件:整条链路唯一在流转的对象 */
@Data
@Builder
public class BinlogEvent {
private String table; // 表名
private OpType opType; // INSERT / UPDATE / DELETE
private long ts; // 变更时间戳
private String primaryKey; // 主键值(用于幂等/去重)
private JSONObject before; // 变更前
private JSONObject after; // 变更后
}
/** 归一化器:每种外部格式一个实现,只做「识别 + 转成领域事件」 */
public interface FormatNormalizer {
boolean support(String rawJson);
BinlogEvent normalize(String rawJson);
}
新增一种采集格式,只需加一个 FormatNormalizer 实现,不改任何下游逻辑。至此原始报文的差异被关在门外——从下一步起,我们讨论的对象永远是干净的 BinlogEvent。
本文的设计主线:很多接入方案会写成一条
监听器 → 基类 → 处理器基类 → 业务处理器的深继承链,把「反序列化、路由、过滤、幂等、业务」全塞进继承体系,方法越长、耦合越深。本文反其道而行——用组合代替继承,把整条链路拆成四个职责单一、可独立测试、可自由拼装的构件,由一个编排者BinlogPipeline串起来:
BinlogPipeline (编排者,唯一入口)
├─ FormatNormalizer 归一化:原始报文 → BinlogEvent
├─ EventRouter 路由:@BinlogTable 注解 → 找到目标 Processor
├─ List<EventFilter> 通用过滤责任链:DDL / 开关 / 幂等,与具体表无关,全表复用
└─ BinlogProcessor 处理:表级/操作级业务过滤(acceptXxx) + 业务动作(onXxx)
下面按 归一化 → 编排 → 路由与过滤 → 业务处理 的顺序逐个拆解。
三、编排者:BinlogPipeline
消息队列的监听器(Kafka @KafkaListener / RocketMQ 等)只做一件事:拿到原始报文,交给 BinlogPipeline。所有真正的逻辑都在 Pipeline 里用组合完成,监听器本身没有任何业务,也不继承任何基类。
@Component
public class BinlogConsumer {
@Resource
private BinlogPipeline pipeline;
/** 批量入口:单条失败不影响其它消息 */
public void onMessages(List<String> messages) {
for (String raw : messages) {
try {
pipeline.accept(raw);
} catch (Exception e) {
// 交由 Pipeline 内部已分级处理,这里只兜最外层
log.error("binlog 消费异常, raw={}", raw, e);
}
}
}
}
BinlogPipeline 是整套设计的核心编排者。它不继承任何东西,而是把四个构件注入进来按顺序组合——这就是「组合优于继承」的落点:每个构件都能单独替换、单独测试,Pipeline 只描述「流程」,不掺杂任何一步的「实现」。
@Component
public class BinlogPipeline {
/** 1. 归一化器集合:屏蔽外部格式差异 */
@Resource
private List<FormatNormalizer> normalizers;
/** 2. 路由:事件 → 目标处理器 */
@Resource
private EventRouter router;
/** 3. 过滤责任链:幂等 / 开关 / 业务过滤,按 @Order 排序注入 */
@Resource
private List<EventFilter> filters;
public void accept(String rawJson) {
// ① 归一化:原始报文 → 标准领域事件
BinlogEvent event = normalize(rawJson);
if (event == null) {
return;
}
// ② 路由:按表名找到唯一处理器;找不到即丢弃(无意义重试)
BinlogProcessor processor = router.route(event.getTable());
if (processor == null) {
log.warn("表[{}]无处理器,跳过", event.getTable());
return;
}
// ③ 过滤责任链:任一环节拒绝则终止,且各 Filter 互不知晓彼此
for (EventFilter filter : filters) {
if (!filter.accept(event)) {
log.info("被[{}]拦截, table={}, pk={}",
filter.getClass().getSimpleName(), event.getTable(), event.getPrimaryKey());
return;
}
}
// ④ 分发到处理器:只按操作类型回调,处理器不感知前三步
processor.handle(event);
}
private BinlogEvent normalize(String rawJson) {
return normalizers.stream()
.filter(n -> n.support(rawJson))
.findFirst()
.map(n -> n.normalize(rawJson))
.orElse(null);
}
}
组合式 Pipeline 带来的几个特性:
| 维度 | 组合式 Pipeline |
|---|---|
| 结构 | Pipeline 持有 4 个平级构件,各自独立 |
| 过滤逻辑 | 通用过滤拆成 EventFilter 责任链,业务过滤内聚各处理器 |
| 扩展方式 | 加一个构件实现即可,Pipeline 零改动 |
| 可测试性 | 每个构件都是普通类,可单独单测 |
normalize 已在第二章说明,下面依次拆解 路由(EventRouter)、过滤链(EventFilter) 和 处理器(BinlogProcessor)。
四、路由、过滤与处理
这一章拆解 Pipeline 里剩下的三个构件。它们彼此独立,谁也不知道谁的存在,全靠 Pipeline 编排。
1、路由:用注解显式声明,而非 Bean 名硬约定
「靠 Bean 名等于表名」是一种脆弱的隐式约定——改个 Bean 名就断了,也无法一个处理器管多张表。这里改用显式注解 @BinlogTable 声明归属,启动时扫描建表:
@Target(ElementType.TYPE)
@Retention(RetentionPolicy.RUNTIME)
public @interface BinlogTable {
String[] value(); // 该处理器负责的表名,可多张
}
@Component
public class EventRouter {
private final Map<String, BinlogProcessor> table2processor = new HashMap<>(16);
/** 启动时扫描所有 BinlogProcessor,按注解声明的表名建立路由表 */
public EventRouter(List<BinlogProcessor> processors) {
for (BinlogProcessor p : processors) {
BinlogTable anno = p.getClass().getAnnotation(BinlogTable.class);
if (anno == null) {
throw new IllegalStateException(p.getClass() + " 缺少 @BinlogTable");
}
for (String table : anno.value()) {
BinlogProcessor exist = table2processor.putIfAbsent(table, p);
if (exist != null) {
throw new IllegalStateException("表[" + table + "]被多个处理器声明");
}
}
}
}
public BinlogProcessor route(String table) {
return table2processor.get(table);
}
}
比起隐式约定,注解路由可读、可校验(重复声明启动即报错)、支持一对多。
2、过滤:责任链,一个 Filter 只管一件事
很多接入方案会把「DDL 过滤、开关控制、幂等防乱序」这几种完全不同的关注点,一股脑塞进一个大而全的前置判断方法里,越写越长、越改越怕。这里把它们拆成独立的 EventFilter,每个只管一件事,用 @Order 控制顺序,Pipeline 依次询问:
public interface EventFilter {
/** 返回 false 表示拦截,事件不再向后流转 */
boolean accept(BinlogEvent event);
}
/** ① DDL / 空数据过滤:结构变更不承载业务数据 */
@Order(10)
@Component
public class DdlFilter implements EventFilter {
@Override
public boolean accept(BinlogEvent event) {
return event.getOpType() != null && event.getAfter() != null;
}
}
/** ② 开关过滤:灰度 / 应急降级时可整体关闭某表消费 */
@Order(20)
@Component
public class SwitchFilter implements EventFilter {
@Resource
private BinlogSwitchConfig switchConfig;
@Override
public boolean accept(BinlogEvent event) {
return switchConfig.isEnabled(event.getTable());
}
}
/** ③ 幂等 + 防乱序:仅当本次 ts 不早于已处理 ts 才放行 */
@Order(30)
@Component
public class IdempotentFilter implements EventFilter {
@Resource
private ProcessedTsStore tsStore; // 记录「主键 → 最近处理时间戳」
@Override
public boolean accept(BinlogEvent event) {
Long last = tsStore.get(event.getTable(), event.getPrimaryKey());
if (last != null && event.getTs() < last) {
return false; // 旧消息,丢弃
}
tsStore.put(event.getTable(), event.getPrimaryKey(), event.getTs());
return true;
}
}
责任链的价值就在于此:DDL、开关、幂等这类逻辑对所有表都通用,做成 Filter 后一次编写、全表复用,不必每接一张表就重抄一遍。注意——责任链只承载这类”与具体表无关”的通用过滤;至于”每张表规则不同、甚至同表增删改各不相同”的业务过滤,不属于这里,它们内聚在各自的处理器中(见下一节)。
3、处理器:表级 / 操作级业务过滤都收在这里
到这一层,前面 EventFilter 责任链只解决了与业务无关的通用过滤(DDL、开关、幂等)。但真正棘手的是另一类过滤:
- 每张表的业务过滤规则都不一样:
order_status只同步某些biz_type,user_account要排除测试账号,stock只关心库存真正变化的行…… - 同一张表,增 / 删 / 改的规则又各不相同:可能”新增全要、更新只在关键字段变化时才要、删除直接忽略”。
这类规则如果继续往 EventFilter 责任链里塞,责任链会迅速膨胀成”每表 × 每操作”一堆 Filter,而且每个 Filter 还得先 if (table.equals("xxx")) 判断是不是自己——这就退化成了大杂烩。正确的归属是:谁的业务谁自己管——把「本表、本操作要不要处理」这件事,内聚到该表处理器内部,和对应的业务动作一一并排。
为此,OpTypeProcessor 为每个操作各配一个前置判断钩子 acceptXxx(默认放行),handle 里「先判断本操作是否接受,再执行本操作」:
public interface BinlogProcessor {
void handle(BinlogEvent event);
}
/** 可选基类:把「每操作的前置判断 + 业务动作」成对固化,非强制继承 */
public abstract class OpTypeProcessor implements BinlogProcessor {
@Override
public final void handle(BinlogEvent event) {
switch (event.getOpType()) {
case INSERT: if (acceptInsert(event)) onInsert(event); break;
case UPDATE: if (acceptUpdate(event)) onUpdate(event); break;
case DELETE: if (acceptDelete(event)) onDelete(event); break;
default: /* ignore */
}
}
// —— 每个操作独立的「要不要处理」判断,默认全放行,子类按表覆写 ——
protected boolean acceptInsert(BinlogEvent event) { return true; }
protected boolean acceptUpdate(BinlogEvent event) { return true; }
protected boolean acceptDelete(BinlogEvent event) { return true; }
// —— 每个操作的业务动作,默认空实现,子类按需覆写 ——
protected void onInsert(BinlogEvent event) {}
protected void onUpdate(BinlogEvent event) {}
protected void onDelete(BinlogEvent event) {}
}
这样一来,一张表的”过滤 + 处理”在自己的类里就是一张清清楚楚的表格:哪个操作、满足什么条件、做什么,增删改互不牵连,也不会污染别的表。
@BinlogTable("order_status")
@Component
public class OrderStatusProcessor extends OpTypeProcessor {
@Resource
private OrderSyncService orderSyncService;
/** 新增:只同步目标业务类型的订单 */
@Override
protected boolean acceptInsert(BinlogEvent event) {
return isTargetBiz(event.getAfter());
}
@Override
protected void onInsert(BinlogEvent event) {
orderSyncService.sync(toOrder(event.getAfter()));
}
/** 更新:目标类型 且 状态确实发生变化才处理(避免无意义写下游) */
@Override
protected boolean acceptUpdate(BinlogEvent event) {
return isTargetBiz(event.getAfter()) && statusChanged(event);
}
@Override
protected void onUpdate(BinlogEvent event) {
orderSyncService.sync(toOrder(event.getAfter()));
}
// 未覆写 acceptDelete → 默认放行;也可覆写为 return false 表示本表不处理删除
@Override
protected void onDelete(BinlogEvent event) {
orderSyncService.remove(event.getPrimaryKey());
}
private boolean isTargetBiz(JSONObject after) {
return "NORMAL".equals(after.getString("biz_type"));
}
private boolean statusChanged(BinlogEvent event) {
return event.getBefore() != null
&& !Objects.equals(event.getBefore().getString("status"),
event.getAfter().getString("status"));
}
private Order toOrder(JSONObject after) {
return after.toJavaObject(Order.class);
}
}
过滤的两层归属,到此彻底分清:
| 过滤类型 | 特征 | 归属 | 例子 |
|---|---|---|---|
| 通用过滤 | 与具体表、具体业务无关,所有表一个样 | EventFilter 责任链,一次编写全表复用 |
DDL、开关、幂等防乱序 |
| 表级业务过滤 | 因表而异 | 该表 Processor 内 |
只同步某 biz_type |
| 操作级业务过滤 | 同一表因增/删/改而异 | 该表 Processor 的 acceptInsert/Update/Delete |
更新时字段无变化则跳过、本表不处理删除 |
换句话说:共性沉到责任链,个性收在处理器,且个性还能细到每个操作。责任链永远保持精简,处理器则完整拥有自己那张表的全部业务判断——这正是”组合”给到的自由度:通用能力可插拔复用,业务规则各自内聚。
4、整体调用时序
消息队列
│ 批量原始报文
▼
BinlogConsumer.onMessages() 最外层兜底
│ 逐条
▼
BinlogPipeline.accept(raw) 编排者,唯一主流程
│
├─► FormatNormalizer.normalize() ① 原始报文 → BinlogEvent
│
├─► EventRouter.route(table) ② 注解路由 → 目标 Processor
│
├─► EventFilter 责任链 ③ 通用过滤:DDL → 开关 → 幂等,任一拒绝即终止
│
└─► BinlogProcessor.handle() ④ 先按操作做表级业务过滤(acceptXxx),通过再写业务(onXxx)
四步各司其职、线性可读,没有任何一层需要「向上追继承链」才能看懂全貌。过滤被清晰地劈成两层——通用过滤在第③步的责任链,表级/操作级业务过滤在第④步的处理器内部,责任链因此永远精简,业务规则则各自内聚、互不干扰。
五、生产落地的关键保障
Binlog 消息天生可能重复、可能乱序,加上下游可能临时故障,以下几件事是正确运行的必要条件。
1、幂等与顺序保证
这正是 IdempotentFilter 存在的意义——它作为责任链的一环,对所有表统一生效:
- 幂等:以
表名 + 主键为键记录最近处理时间戳,重复消息被拦截。 - 顺序:比较本次
event.getTs()与已记录时间戳,旧消息直接丢弃(见第四章IdempotentFilter)。 - 下游若能保证幂等(如
INSERT ... ON DUPLICATE KEY UPDATE),则是双重保险。
因为它是独立 Filter,业务处理器完全不用重复写这段逻辑。
2、异常与重试分级
不要所有异常都无脑重试。可在 BinlogConsumer / BinlogPipeline 外层统一按异常类型分级:
| 异常类型 | 处理策略 |
|---|---|
| 归一化 / 反序列化失败 | 记录告警 + 丢弃(重试无意义) |
| 下游临时故障 | 向上抛出 → 消息队列重试 |
| 重试超限 | 投递到死信队列人工介入 |
分级逻辑集中一处,各构件只管抛出语义清晰的异常,不各自决定重试策略。
3、可观测性与资源隔离
- 日志:在 Pipeline 统一打印
table + opType + 主键 + ts,全链路一眼可查。 - 监控:采集消费 TPS、失败率、消费延迟(
now - ts)并接入告警。 - 线程池隔离:按表路由后,不同表可投递到独立线程池处理,避免慢表拖垮快表。
4、接入新表的标准步骤(SOP)
- 定义该表对应的 BO 对象(字段与
after一致)。 - 新增 处理器:
@BinlogTable("表名")+ 继承OpTypeProcessor,实现onInsert/onUpdate/onDelete。 - 如是新采集格式,补充一个
FormatNormalizer实现。 - 如该表/某个操作有业务过滤规则,在处理器里覆写对应的
acceptInsert/acceptUpdate/acceptDelete即可;配置该表 Topic 订阅与开关。(只有当过滤规则真的对所有表通用时,才考虑新增EventFilter)
全程不改任何已有代码,只做”加法”——这正是组合式设计的红利。
六、小结
| 环节 | 构件 | 职责 |
|---|---|---|
| 来源 | MySQL Binlog + 采集组件 | 捕获变更并投递消息队列 |
| 归一化 | FormatNormalizer |
外部报文 → 标准 BinlogEvent |
| 编排 | BinlogPipeline |
组合四个构件,描述主流程 |
| 路由 | EventRouter + @BinlogTable |
显式注解路由,启动即校验 |
| 通用过滤 | EventFilter 责任链 |
DDL / 开关 / 幂等,与表无关,全表复用 |
| 业务过滤+处理 | BinlogProcessor |
表级/操作级过滤(acceptXxx) + 业务动作(onXxx) |
核心思想:用组合代替继承,把「接入」拆成一组职责单一、可独立测试、可自由拼装的构件。过滤被清晰劈成两层——通用能力(DDL/开关/幂等)沉淀为可复用的 EventFilter 责任链,因表/因操作而异的业务过滤则内聚在各自处理器的 acceptXxx 里,责任链永远精简,业务规则各自独立;接入新表永远是”加一个构件”,而非”改一条继承链”。

