如何在Debezium PostgreSQL连接器的after对象中添加静态source字段?
问题描述
使用Debezium PostgreSQL源连接器实现Kafka CDC功能时,需要在推送至Kafka的消息after对象中添加一个静态不变的source字段(值为C),该字段不存在于源数据库,且无法修改数据库或创建视图。当前配置错误地在消息根节点添加了after.source字段,而非嵌入到after对象内部。
期望的消息Schema
{ "before": null, "after": { "rid": "3b99c447-65a8-4d6b-bbff-2c33b7944696", "cust": 75862, "loc": 916719, "meter": "A90OC5385040", "cosum": "2.06", "cosdt": 1673330400000000, "costy": "I", "source": "C" }, "source": { "version": "2.4.2.Final", "connector": "postgresql", "name": "mmv2_pgami", "ts_ms": 1711944632077, "snapshot": "false", "db": "Pgami_db", "sequence": "[null,\"23516504\"]", "schema": "public", "table": "mreads", "txId": 574, "lsn": 23516504, "xmin": null }, "op": "c", "ts_ms": 1711944632426, "transaction": null }
当前错误配置
{ "name": "postgres_connector", "config": { "connector.class": "io.debezium.connector.postgresql.PostgresConnector", // 其他配置项... // 错误的AddFields转换配置 "transforms": "addSourceField", "transforms.addSourceField.type": "org.apache.kafka.connect.transforms.InsertField$Value", "transforms.addSourceField.static.field": "after.source", "transforms.addSourceField.static.value": "C" } }
正确配置方案
InsertField转换不支持直接通过嵌套路径添加字段,需要通过三步转换实现:先将after字段提升为顶层结构,插入静态字段后,再将其重新包装回原有的after字段中。
完整配置示例
{ "name": "postgres_connector", "config": { "connector.class": "io.debezium.connector.postgresql.PostgresConnector", // 其他必要配置(如数据库连接、表过滤等)... "transforms": "hoistAfter,addSource,wrapAfter", // 第一步:将after字段提升为顶层值 "transforms.hoistAfter.type": "org.apache.kafka.connect.transforms.HoistField$Value", "transforms.hoistAfter.field": "after", // 第二步:插入静态source字段 "transforms.addSource.type": "org.apache.kafka.connect.transforms.InsertField$Value", "transforms.addSource.static.field": "source", "transforms.addSource.static.value": "C", // 第三步:将修改后的结构重新包装回after字段 "transforms.wrapAfter.type": "org.apache.kafka.connect.transforms.WrapField$Value", "transforms.wrapAfter.field": "after" } }
配置说明
- hoistAfter:用
HoistField把消息里的after对象提取为顶层值,后续转换将直接作用于这个对象。 - addSource:用
InsertField在当前顶层对象(即原after)中添加静态source字段,值设为C。 - wrapAfter:用
WrapField把修改后的对象重新包装回after字段,恢复原消息结构。
按此配置后,生成的Kafka消息会在after对象内部包含指定的静态source字段,符合预期的Schema结构。
内容的提问来源于stack exchange,提问作者user1913559
相关产品推荐
相关产品推荐

