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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 07:57:38