Debezium ExtractNewRecordState转换失效问题排查求助
解决Debezium ExtractNewRecordState转换未生效的问题
你这明显是踩了个容易忽略的拼写坑!你的配置里用了单数的transform,但Kafka Connect要求转换配置的参数名必须是复数的transforms——这就是转换完全没起作用的核心原因。
问题根源
你当前的配置片段:
{ // 省略其他配置 "transform": "unwrap", "transform.unwrap.type": "io.debezium.transforms.ExtractNewRecordState" }
这里的transform是单数写法,不符合Kafka Connect的配置规范,导致Connect直接忽略了你的转换规则,所以消息还是以Debezium原始的CDC格式(包含before、source等所有字段)发送到Kafka Topic中。从你提供的日志里反复出现transforms = []也能佐证这一点——Connect根本没加载到你的转换配置。
正确的配置示例
把参数名改成复数,同时可以补充一些实用的转换配置(比如处理删除事件、丢弃无用的墓碑消息):
{ // 其他源连接器配置... "transforms": "unwrap", "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState", "transforms.unwrap.delete.handling.mode": "rewrite", "transforms.unwrap.drop.tombstones": "true" }
验证步骤
- 更新连接器配置(可以用Kafka Connect的REST API重新提交,或者重启连接器)
- 再次消费Kafka Topic的消息,你应该会看到消息只保留
after字段里的业务数据,比如:
{"id":1,"name":"ggg"}
- 查看Connect日志,此时应该不会再出现
transforms = []的空数组记录,而是能看到正确加载unwrap转换的日志内容。
另外提醒一句:这个转换只需要在源连接器端配置就够了,不用同时在源和Sink两端重复配置,这样Sink收到的就是处理后的干净数据,更高效。
内容的提问来源于stack exchange,提问作者user8510613
相关产品推荐
相关产品推荐

