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

