MongoDB数据迁移时_id格式不一致的修正方案咨询
MongoDB数据迁移:统一_id字段格式方案
问题说明
跨MongoDB实例迁移数据时,出现_id格式不匹配:
- 目标端要求的正确格式:
"_id": { "$oid": "57f2773113b28328a0dabfe2" }
- 当前源端输出的错误格式(字符串化的ObjectID):
_id: "{\"$oid\": \"57f2773113b28328a0dabfe4\"}"
解决办法
方法1:修改源连接器配置,直接输出规范格式
你的MongoDB Atlas源连接器当前使用DefaultJson格式输出,会把BSON类型转为字符串。改成ExtendedJson就能输出符合MongoDB规范的嵌套$oid结构:
在源连接器的config中修改:
"output.json.format": "ExtendedJson"
保留output.data.format": "JSON"不变,重启源连接器即可。
方法2:通过Kafka Connect转换解析字符串化的_id
如果源端配置无法调整,就在链路中加转换逻辑,把字符串转成嵌套对象。比如给Sink连接器加ScriptTransform:
在Sink连接器的config中添加:
"transforms": "parseId", "transforms.parseId.type": "org.apache.kafka.connect.transforms.Script", "transforms.parseId.script": "value._id = JSON.parse(value._id); return value;"
注意:需要Kafka Connect环境支持
ScriptTransform,提前安装对应插件。
方法3:调整Sink连接器的ID策略
当前Sink用的FullKeyStrategy会把消息Key作为文档ID。如果源端已经输出正确的_id结构,改成ProvidedInValueStrategy让Sink直接用Value里的_id:
在Sink连接器的config中修改:
"doc.id.strategy": "ProvidedInValueStrategy", "doc.id.strategy.overwrite.existing": "true"
验证
- 改完配置重启连接器
- 发测试数据到Kafka,检查目标库的_id格式
- 全量迁移的话,确认
copy.existing逻辑也能正确转换格式
内容的提问来源于stack exchange,提问作者sai jyothsna pentyala
相关产品推荐
相关产品推荐

