如何将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
相关产品推荐
相关产品推荐

