如何通过Kafka Connect Transforms遮蔽或删除Debezium同步的JSON列
解决方案
配置不生效原因
你的原有配置失效通常是两个原因导致:
- MaskField仅对字符串、数字等基础类型有默认替换值,JSON/Struct类型的复杂字段需要显式指定
replacement参数才能生效 - Debezium默认输出的消息是包含
before/after/source等元字段的Envelope封装结构,业务字段实际嵌套在after字段下,直接指定顶级data字段无法匹配到目标字段。
方案1:删除data字段(对应期望格式2)
推荐使用ReplaceField转换器实现,配置如下:
已开启Debezium解包(业务字段在消息顶级)
transforms: "removeDataField" transforms.removeDataField.type: "org.apache.kafka.connect.transforms.ReplaceField$Value" # Kafka 3.0+版本用exclude,低版本用blacklist transforms.removeDataField.exclude: "data"
未开启Debezium解包
需要先加Envelope解包转换,再删除字段:
transforms: "unwrap,removeDataField" # 第一步:解包Debezium封装结构,提取变更后的业务数据 transforms.unwrap.type: "io.debezium.transforms.ExtractNewRecordState" transforms.unwrap.drop.tombstones: "true" # 第二步:删除data字段 transforms.removeDataField.type: "org.apache.kafka.connect.transforms.ReplaceField$Value" transforms.removeDataField.exclude: "data"
方案2:将data字段置为null(对应期望格式1)
使用MaskField转换器,显式指定替换值为null即可:
transforms: "unwrap,maskDataField" transforms.unwrap.type: "io.debezium.transforms.ExtractNewRecordState" transforms.unwrap.drop.tombstones: "true" transforms.maskDataField.type: "org.apache.kafka.connect.transforms.MaskField$Value" transforms.maskDataField.fields: "data" transforms.maskDataField.replacement: null
注意:如果开启了Schema校验,需要保证
data字段在Schema中是可选属性(optional=true),否则会报错。如果data是必填字段,可以先加SetSchemaMetadata转换将其改为可选;如果配置了value.converter.schemas.enable=false关闭了Schema,无需处理校验问题直接配置即可生效。
内容的提问来源于stack exchange,提问作者user8606365
相关产品推荐
相关产品推荐

