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

如何在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"
    }
}

配置说明

  1. hoistAfter:用HoistField把消息里的after对象提取为顶层值,后续转换将直接作用于这个对象。
  2. addSource:用InsertField在当前顶层对象(即原after)中添加静态source字段,值设为C。
  3. wrapAfter:用WrapField把修改后的对象重新包装回after字段,恢复原消息结构。

按此配置后,生成的Kafka消息会在after对象内部包含指定的静态source字段,符合预期的Schema结构。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 16:33:13