Debezium作为一款优秀的CDC数据变更捕获工具,能够实时捕获MySQL、PostgreSQL、MongoDB等数据库的binlog或WAL日志,并将变更以事件的形式写入Kafka。不过Debezium产生的事件结构是相当“厚”的:一条消息里包含了before(变更前的值)、after(变更后的值)、source(源库元信息)、op(操作类型)等多个嵌套层级。对于很多下游消费者来说,尤其是做数据同步到Elasticsearch、Redis或者数据湖的场景,这种嵌套结构反而增加了消费端的解析负担。这时候就需要用到Kafka Connect的SMT机制,把事件结构在写入Kafka之前直接扁平化成普通的一行记录。

一、先看懂Debezium的原始事件结构
在没有做任何转换之前,Debezium发送到Kafka的一条更新事件大概长成下面这个样子:
{
"before": {
"id": 1001,
"user_name": "张三",
"balance": 500
},
"after": {
"id": 1001,
"user_name": "张三",
"balance": 800
},
"source": {
"db": "order_db",
"table": "account",
"lsn": 123456,
"ts_ms": 1700000000000
},
"op": "u",
"ts_ms": 1700000000500
}
可以看到,真正的业务数据被包裹在after字段里面,消费者如果只想要最新的字段值,就必须先取出after再做一层解析。而且before和after的结构完全一样,一些下游存储(比如Elasticsearch的动态mapping)会因为字段类型不确定而报错。这就是为什么很多团队在采集链路中第一时间要做的优化就是事件扁平化。
另一个需要理解的点是,Debezium的连接器本质上是Kafka Connect的一个Source Connector,Kafka Connect框架允许在连接器配置中声明一组转换器(transforms),这些转换器会在消息进入Kafka Topic之前依次执行。理解了这个执行时机,就能明白为什么SMT属于“前置处理”,它发生在序列化之前,效率远高于消费者端再解析一次。
二、ExtractNewRecordState转换器的配置方法
扁平化事件的核心是Debezium官方提供的io.debezium.transforms.ExtractNewRecordState这个SMT,它的作用是把嵌套的after对象提升为消息的顶层字段。完整的连接器配置示例如下:
{
"name": "mysql-connector",
"config": {
"connector.class": "io.debezium.connector.mysql.MySqlConnector",
"database.hostname": "192.168.0.1",
"database.port": "3306",
"database.user": "debezium",
"database.password": "password",
"database.server.id": "184055",
"topic.prefix": "order_server",
"table.include.list": "order_db.account",
"transforms": "unwrap",
"transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
"transforms.unwrap.drop.tombstones": "false",
"transforms.unwrap.delete.handling.mode": "rewrite",
"transforms.unwrap.add.fields": "op,table,lsn",
"transforms.unwrap.add.headers": "db"
}
}
配置中的transforms声明了转换器的别名,这里叫unwrap,后面的所有参数都以transforms.unwrap.为前缀。经过这一层转换后,写入Kafka的消息就变成了纯粹的键值结构:
{
"id": 1001,
"user_name": "张三",
"balance": 800,
"__op": "u",
"__table": "account",
"__lsn": "123456"
}
业务字段直接暴露在顶层,op、table、lsn等元信息通过add.fields参数以双下划线前缀附加进来,消费者既拿到了干净的数据,又保留了必要的溯源信息。如果不希望字段名带前缀,还可以使用add.fields.prefix参数自定义前缀,或者设为空字符串去掉前缀。
三、删除事件与墓碑消息的处理策略
扁平化最容易被忽略的坑就是删除事件。Debezium在捕获到DELETE操作时会发送一条值为null的删除事件,随后再发送一条墓碑消息(Tombstone)用于Kafka日志压缩。如果在SMT中简单使用默认配置,消费者收到的删除消息可能是空值,根本不知道哪条数据被删了。
这时候delete.handling.mode参数就非常关键了。它有三个可选值:none是默认行为,直接丢弃删除事件的有效信息;rewrite会保留被删除记录的主键值,并额外添加一个__deleted字段标记为true,这样下游只需要按主键同步删除即可;drop则是干脆不发送删除相关的消息。对于同步到Elasticsearch这类需要显式执行删除操作的场景,推荐使用rewrite模式,消费者读取到__deleted为true的记录时主动调用delete API。
{
"id": 1001,
"__op": "d",
"__deleted": "true"
}
关于drop.tombstones参数,如果下游消费框架对null消息不友好,可以设置为true把墓碑消息也过滤掉。但要注意,如果你的Topic开启了日志压缩(log.cleanup.policy=compact),墓碑消息是必要的删除标记,随意丢弃会导致旧数据无法被清理。
四、扁平化的代价:什么场景不适合拍平
事件扁平化虽然让消费端省事,但它是有信息损失的。最直接的一点是before前镜像值被丢弃了。如果你的业务需要审计功能,比如记录“余额从500变成800”这样的变更轨迹,或者下游需要基于旧值做条件判断(例如只有库存从正数变为零时才触发告警),拍平之后的消息就无法满足需求,此时应该保留完整的Debezium事件结构,在消费端自行解析。
其次,扁平化之后消息结构因表而异,多表共用一个Topic的场景下(通过RegexRouting路由),不同表的字段混合在一起可能产生字段名冲突。解决办法是配合ExtractNewRecordState的field命名前缀,或者在路由阶段就按表拆分Topic。对于 MongoDB 连接器,由于文档结构本身就是嵌套的,还需要额外关注嵌套数组的处理,必要时使用add.fields补充collection信息来区分来源。
总结来看,扁平化适合“数据搬运”类场景,目标是把最新状态同步到缓存、搜索索引或宽表中;而保留完整事件结构适合“事件驱动”类场景,需要感知每一次变更的细节。两种方案没有绝对优劣,关键是根据下游消费方式做出取舍。配置时建议先在测试环境验证删除事件和边界情况,确认消息格式符合预期后再上生产环境。