You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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"
}

验证步骤

  1. 更新连接器配置(可以用Kafka Connect的REST API重新提交,或者重启连接器)
  2. 再次消费Kafka Topic的消息,你应该会看到消息只保留after字段里的业务数据,比如:
{"id":1,"name":"ggg"}
  1. 查看Connect日志,此时应该不会再出现transforms = []的空数组记录,而是能看到正确加载unwrap转换的日志内容。

另外提醒一句:这个转换只需要在源连接器端配置就够了,不用同时在源和Sink两端重复配置,这样Sink收到的就是处理后的干净数据,更高效。

内容的提问来源于stack exchange,提问作者user8510613

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.07 16:47:34