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

Debezium MongoDB使用ExtractNewDocumentState SMT添加非元数据header失败

解决方案

根因说明

MongoDB 连接器的 ExtractNewDocumentState SMT 输出默认是 JSON 字符串格式,而非 Postgres 连接器 ExtractNewRecordState 输出的 Struct 结构化对象,同时该 SMT 执行后已经将变更后的文档内容(原after节点内容)作为记录的核心 value 输出,无需再携带after.前缀寻址,两种原因叠加导致你触发字段不存在的报错。


可用解决方案

方案1:调整内置SMT配置(无额外依赖,推荐)

修改ExtractNewDocumentState的输出格式为 Struct 类型,同时去掉add.headers配置中的after.前缀即可,调整后完整配置如下:

"transforms": "unwrap",
"transforms.unwrap.type":"io.debezium.connector.mongodb.transforms.ExtractNewDocumentState",
"transforms.unwrap.output.type": "struct",
"transforms.unwrap.add.headers": "Id:Id",
"transforms.unwrap.add.headers.prefix": ""

该方案无需引入额外组件,直接使用Debezium内置能力即可实现需求。

方案2:串联JSON解析SMT(适配需要保留JSON字符串输出的场景)

如果你业务需要保留ExtractNewDocumentState输出为JSON字符串的默认行为,可以串联通用JSON字段提取SMT实现需求,配置示例如下:

"transforms": "unwrap,extractId",
"transforms.unwrap.type":"io.debezium.connector.mongodb.transforms.ExtractNewDocumentState",
"transforms.extractId.type": "com.github.jcustenborder.kafka.connect.transform.common.ExtractFromJSON$Value",
"transforms.extractId.header.name": "Id",
"transforms.extractId.json.path": "$.Id"

方案3:自定义SMT(无第三方依赖场景)

如果不允许引入第三方SMT,可以自定义开发单功能SMT,核心逻辑为:读取记录的value(JSON字符串)、解析后提取指定字段值、写入到消息header,仅需几十行代码即可实现。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 19:15:02