MongoDB Kafka Source Connector Pipeline触发BsonInvalidOperationException故障排查
解决MongoDB Kafka Source Connector Pipeline转换异常问题
问题场景
需要通过MongoDB Kafka Source Connector将集合文档转换为指定时序结构:
- 源集合文档结构:
{ "_id": "6798680c885bb2a8fbe93485", "_ts": { "$date": "2025-01-28T05:15:56.507Z" }, "bId": 226, "fingprint": "4d455166-2166-4670-9315-8d9b7c7024a7", "cId": 682, "timestamp": { "$numberLong": "1738041356" }, "catId": 1 }
- 期望转换后的结构:
{ "_ts": { "$date": "2025-01-28T05:15:56.507Z" }, "metafield": { "bId": 226, "cId": 682, "catId": 1 }, "_id": "6798680c885bb2a8fbe93485", "fingprint": "4d455166-2166-4670-9315-8d9b7c7024a7", "timestamp": { "$numberLong": "1738041356" } }
配置带Pipeline的Kafka源连接器后注册成功,但触发异常:
[2025-04-24 07:43:11,032] ERROR WorkerSourceTask{id=mongo-source-velocity-0} Task threw an uncaught and unrecoverable exception. Task is being killed and will not recover until manually restarted (org.apache.kafka.connect.runtime.WorkerTask) org.bson.BsonInvalidOperationException: Value expected to be of type DOCUMENT is of unexpected type STRING at org.bson.BsonValue.throwIfInvalidType(BsonValue.java:419) at org.bson.BsonValue.asDocument(BsonValue.java:47) at org.bson.BsonDocument.getDocument(BsonDocument.java:154) at com.mongodb.kafka.connect.source.StartedMongoSourceTask.pollInternal(StartedMongoSourceTask.java:217) at com.mongodb.kafka.connect.source.StartedMongoSourceTask.poll(StartedMongoSourceTask.java:190) at com.mongodb.kafka.connect.source.MongoSourceTask.poll(MongoSourceTask.java:180) at org.apache.kafka.connect.runtime.WorkerSourceTask.poll(WorkerSourceTask.java:305) at org.apache.kafka.connect.runtime.WorkerSourceTask.execute(WorkerSourceTask.java:249) at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:188) at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:243) at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515) at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) at java.base/java.lang.Thread.run(Thread.java:829)
未配置Pipeline时连接器可正常复制数据,但需要在发布到Kafka前完成结构转换。
错误原因分析
配置中设置了publish.full.document.only=true,此时Pipeline直接作用于集合的原始文档,不需要使用$fullDocument前缀来引用字段。原Pipeline中使用$fullDocument.bId这类写法,会导致MongoDB将其解析为字符串类型,而非预期的文档字段引用,触发类型不匹配异常。
修正后的连接器配置
修改Pipeline部分,移除所有$fullDocument.前缀,直接引用字段名:
{ "name": "mongo-source-data", "config": { "connector.class": "com.mongodb.kafka.connect.MongoSourceConnector", "key.converter": "org.apache.kafka.connect.storage.StringConverter", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "publish.full.document.only": "true", "connection.uri": "mongodb://mongo1", "topic.prefix": "data", "database": "test_db", "collection": "test_cn", "copy.existing": "true", "pipeline": "[{\"$project\":{\"metafield\":{\"bId\":\"$bId\",\"cId\":\"$cId\",\"catId\":\"$catId\"},\"_ts\":\"$_ts\",\"fingprint\":\"$fingprint\",\"timestamp\":\"$timestamp\",\"_id\":\"$_id\"}}]" } }
验证说明
修正Pipeline后,连接器会正确解析字段引用,将源文档转换为期望的时序结构并发布到Kafka,同时避免类型异常。
内容的提问来源于stack exchange,提问作者Zeeshan
相关产品推荐
相关产品推荐

