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

如何将Debezium生成的Kafka消息转换为指定JSON格式?

使用Kafka Sink Transforms转换Debezium消息格式

Debezium生成的消息是包含schema和payload的Envelope结构,要转换成仅含业务字段和__op操作标识的JSON格式,可通过Kafka Connect的内置Transforms组合实现,以下是具体配置和说明:

核心配置示例

假设你的Sink Connector配置如下(以JDBC Sink为例,其他Sink可复用Transform部分):

name=business-data-sink
connector.class=io.confluent.connect.jdbc.JdbcSinkConnector
topics=your-debezium-source-topic
connection.url=jdbc:mysql://your-db-host:3306/your-db
connection.user=db-user
connection.password=db-pass
auto.create=true
auto.evolve=true

# 配置Transforms,按顺序执行
transforms=extractPayload,flattenBizFields,renameOp,cleanupFields

# 1. 提取payload字段,丢弃外层schema结构
transforms.extractPayload.type=org.apache.kafka.connect.transforms.ExtractField$Value
transforms.extractPayload.field=payload

# 2. 展开after中的业务字段到顶层(替换嵌套结构)
transforms.flattenBizFields.type=org.apache.kafka.connect.transforms.Flatten$Value
transforms.flattenBizFields.delimiter=  # 留空直接移除after前缀,如after.id → id

# 3. 将Debezium原生op字段重命名为__op
transforms.renameOp.type=org.apache.kafka.connect.transforms.ReplaceField$Value
transforms.renameOp.renames=op:__op

# 4. 过滤掉不需要的Debezium元数据字段
transforms.cleanupFields.type=org.apache.kafka.connect.transforms.ReplaceField$Value
transforms.cleanupFields.blacklist=before,source,ts_ms,transaction

各Transform作用说明

  • ExtractField$Value:从原始消息中剥离出payload部分作为新的消息体,丢弃冗余的schema字段。
  • Flatten$Value:将payload.after下的嵌套业务字段提升到消息顶层,消除层级结构,符合业务系统的字段格式要求。
  • ReplaceField$Value(重命名):将Debezium原生的操作标识字段op(取值为c=创建、u=更新、d=删除、r=快照)重命名为业务需要的__op。
  • ReplaceField$Value(过滤):通过黑名单移除before(更新/删除前的旧数据)、source(源端元数据)等无关字段,只保留业务字段和__op。

特殊场景处理(可选)

如果需要处理删除操作(op=d时after字段为null),可添加分支Transform分别处理新增/更新和删除场景:

transforms=branchRecords,extractPayloadUpsert,extractPayloadDelete,flattenUpsert,flattenDelete,renameOpUpsert,renameOpDelete

# 按op值拆分消息到不同分支
transforms.branchRecords.type=org.apache.kafka.connect.transforms.Branch$Value
transforms.branchRecords.if=value.payload.op == 'd'
transforms.branchRecords.branches=deleteBranch,upsertBranch
transforms.branchRecords.topic.suffix.deleteBranch=-delete
transforms.branchRecords.topic.suffix.upsertBranch=-upsert

# 处理新增/更新分支:提取after
transforms.extractPayloadUpsert.type=org.apache.kafka.connect.transforms.ExtractField$Value
transforms.extractPayloadUpsert.field=payload.after
transforms.extractPayloadUpsert.transforms=upsertBranch

# 处理删除分支:提取before(作为删除依据)
transforms.extractPayloadDelete.type=org.apache.kafka.connect.transforms.ExtractField$Value
transforms.extractPayloadDelete.field=payload.before
transforms.extractPayloadDelete.transforms=deleteBranch

# 后续的flatten、renameOp配置分别对应两个分支即可

内容的提问来源于stack exchange,提问作者Albert T. Wong

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 06:52:03