debezium要求mysql必须开启binlog且配置为row格式、binlog_row_image=full、server_id唯一非零;普通账号需授予replication slave、replication client和select权限;connector配置中database.server.name和table.include.list为必填项,snapshot.mode推荐initial。

MySQL Binlog 必须开启且配置正确
Debezium 依赖 MySQL 的 binlog 实时读取变更,如果 binlog_format 不是 ROW,或 binlog_row_image 不是 FULL,会直接丢变更或报错 Unsupported binlog format。另外,server_id 必须全局唯一且非零,否则 Kafka Connect 启动后反复重试连接。
-
my.cnf中至少需包含:[mysqld] server_id = 18273 log_bin = mysql-bin binlog_format = ROW binlog_row_image = FULL expire_logs_days = 7
- 执行
SHOW VARIABLES LIKE 'binlog_format';和SHOW VARIABLES LIKE 'binlog_row_image';确认生效,修改后必须重启 MySQL - 普通用户账号需授予
REPLICATION SLAVE, REPLICATION CLIENT, SELECT权限,仅用SELECT会卡在初始化阶段,日志报Access denied; you need (at least one of) the SUPER, REPLICATION CLIENT privilege(s)
Debezium MySQL Connector 配置关键项不能漏
Connector JSON 配置里,database.server.name 是逻辑名,后续 Kafka topic 名称和 schema registry 中的命名都基于它;table.include.list 必须显式指定,空值或通配符(如 db1.*)在新版 Debezium(2.0+)中默认被禁用,不填就收不到任何事件。
- 最小可用配置示例:
{ "name": "mysql-connector", "config": { "connector.class": "io.debezium.connector.mysql.MySqlConnector", "tasks.max": "1", "database.hostname": "mysql-host", "database.port": "3306", "database.user": "debezium", "database.password": "xxx", "database.server.id": "18273", "database.server.name": "mysql-server-1", "database.include.list": "inventory", "table.include.list": "inventory.customers,inventory.orders", "snapshot.mode": "initial" } } -
snapshot.mode推荐用initial(首次全量 + 增量),生产环境避免用never,否则丢失历史数据;若表大,可配合snapshot.fetch.size控制单次拉取行数 - 不要设
database.history.kafka.bootstrap.servers指向集群外地址——该配置只用于内部 offset 存储,必须和 Kafka Connect worker 使用同一套 broker
消费端解析 Debezium Avro 事件要注意结构嵌套
Debezium 输出的每条消息 value 是 Avro 格式,顶层包含 before、after、source、op、ts_ms 字段。直接用 JSON.parse() 或简单反序列化会失败,因为实际 payload 被包装在 after(INSERT/UPDATE)或 before(DELETE)里,且字段名带双下划线前缀(如 __deleted)。
- 典型变更消息结构:
{ "schema": { ... }, "payload": { "before": null, "after": {"id": 1001, "name": "Alice"}, "source": {"version":"2.4.0.Final", "name":"mysql-server-1", "table":"customers", ...}, "op": "c", "ts_ms": 1712345678901 } } - 真正业务数据在
payload.after(新增/更新)或payload.before(删除),op值为c(create)、u(update)、d(delete)、r(read,快照期间) - 若用 Kafka Streams 或 Flink 处理,需注册对应 Avro schema(来自 Schema Registry),不能靠自动推断;本地调试可用
kafka-avro-console-consumer加--property print.key=true --property key.separator=" : "查看原始结构
跨系统写入时主键冲突和幂等性必须自己兜底
Debezium 只负责“发出变更”,不保证下游系统写入成功或去重。比如目标库已存在同主键记录,或网络重试导致同一条 op=c 消息被消费两次,就会写入重复数据。
- 推荐在消费端做轻量级幂等:用
payload.source.ts_ms+payload.source.table+payload.after.id组合成唯一 key,缓存最近 N 分钟的 key 到 Redis,重复则跳过 - 对 UPDATE 场景,别直接
INSERT ... ON DUPLICATE KEY UPDATE,因为 Debezium 的after不含完整字段(比如未更新的字段为 null),应先查再合并,或用 CDC-aware 的目标端(如 Apache Doris 的REPLACE表模型) - DELETE 操作尤其危险:一旦消费延迟,又遇到目标库自动清理机制(如 TTL),可能误删本不该删的数据;建议改为软删字段标记 + 定期归档,而非物理删除
真正难的不是把变更发出来,而是让下游系统能按 MySQL 的事务边界、顺序和语义还原状态——这需要消费端对 op 类型、ts_ms、source.event_id 做联合判断,而不是当成普通消息流处理。











