MongoDB Kafka Connect Sink连接器更新操作同步失败
问题根因
本质是版本不兼容:
truncatedArrays是MongoDB 6.0+版本变更流返回的updateDescription里新增的字段,用来记录更新操作中被截断的数组内容- 你当前用的Mongo Kafka Sink Connector版本偏老,内置的
ChangeStreamHandler没适配这个新字段,代码逻辑里只要发现updateDescription里除了updatedFields、removedFields之外有其他字段,就直接抛异常阻断流程,避免漏数据 - Insert操作不需要解析
updateDescription节点,所以同步正常。
解决方案
按落地成本从低到高排列:
方案1:升级Mongo Kafka Connect插件版本
最省事的根治方法,把连接器升级到1.8.0及以上正式版即可。官方从1.8.0版本开始已经适配了MongoDB 6.0的truncatedArrays字段,升级完不用修改现有配置,重启连接器就能正常处理update事件。
注意:升级前先在测试环境跑通全量+增量同步逻辑,避免版本差异带来其他兼容问题。
方案2:源端管道直接剔除truncatedArrays字段
如果暂时没法升级连接器,直接在Source端的聚合管道里加投影步骤,把truncatedArrays字段过滤掉,不让它传到Sink端即可。修改后的Source端pipeline配置参考:
"pipeline": "[ {\"$match\":{\"operationType\": { \"$in\": [ \"update\",\"insert\" ]}}}, {\"$project\": { \"_id\": 1, \"operationType\": 1, \"clusterTime\": 1, \"ns\": 1, \"documentKey\": 1, \"updateDescription.updatedFields\": 1, \"updateDescription.removedFields\": 1, \"fullDocument\": 1 } } ]"
这个方案存在局限性:如果业务后续会用$push搭配$slice之类的操作做数组截断,目标端不会同步这部分变更,仅适合完全没有数组截断操作的场景。
方案3:自定义CDC处理逻辑
如果既无法升级版本,又需要支持数组截断的同步,可以自己实现类继承官方的ChangeStreamHandler,重写update事件的解析逻辑,兼容truncatedArrays字段的处理,打包后放到Connect的插件目录,再把Sink端的change.data.capture.handler配置改成自定义类的全限定名即可。这个方案开发维护成本最高,无特殊需求不推荐选用。
验证步骤
修改配置后按以下顺序验证:
- 重启Source、Sink连接器,检查启动日志无异常
- 往源集合
transactionV2插入一条测试数据,确认目标集合transactionV1正常接收 - 对源集合的测试文档执行一次普通update操作(比如修改amount字段),确认目标端对应文档同步更新,日志中无之前的字段非预期异常
- 如果选用方案2,额外测试数组截断类更新,确认同步结果符合业务预期
内容的提问来源于stack exchange,提问作者Ambuj Mehra
相关产品推荐
相关产品推荐

