使用Kafka MongoDB Sink Connector时$oid字段报错求助
问题:Kafka MongoDB Sink Connector因$oid字段报错无法写入文档
使用Kafka MongoDB Sink Connector推送Kafka Topic中的JSON文档到MongoDB时,因文档包含$oid字段触发报错,无法完成写入。
错误信息
{"name":"mongodb-sink-connector","connector":{"state":"RUNNING","worker_id":"localhost:8083"},"tasks":[{"id":0,"state":"FAILED","worker_id":"localhost:8083","trace":"org.apache.kafka.connect.errors.ConnectException: Exiting WorkerSinkTask due to unrecoverable exception.\n at org.apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:610)\n at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:330)\n at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:232)\n at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:201)\n at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:188)\n at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:237)\n at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515)\n at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264)\n at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128)\n at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628)\n at java.base/java.lang.Thread.run(Thread.java:829)\nCaused by: org.apache.kafka.connect.errors.DataException: Failed to write mongodb documents\n at com.mongodb.kafka.connect.sink.MongoSinkTask.bulkWriteBatch(MongoSinkTask.java:227)\n at java.base/java.util.ArrayList.forEach(ArrayList.java:1541)\n at com.mongodb.kafka.connect.sink.MongoSinkTask.put(MongoSinkTask.java:122)\n at org.apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:582)\n ... 10 more\nCaused by: java.lang.IllegalArgumentException: Invalid BSON field name $oid\n at org.bson.AbstractBsonWriter.writeName(AbstractBsonWriter.java:534)\n at com.mongodb.internal.connection.BsonWriterDecorator.writeName(BsonWriterDecorator.java:193)\n at org.bson.codecs.BsonDocumentCodec.encode(BsonDocumentCodec.java:117)\n at org.bson.codecs.BsonDocumentCodec.encode(BsonDocumentCodec.java:42)\n at org.bson.codecs.EncoderContext.encodeWithChildContext(EncoderContext.java:91)\n at org.bson.codecs.BsonDocumentCodec.writeValue(BsonDocumentCodec.java:139)\n at org.bson.codecs.BsonDocumentCodec.encode(BsonDocumentCodec.java:118)\n at org.bson.codecs.BsonDocumentCodec.encode(BsonDocumentCodec.java:42)\n at com.mongodb.internal.connection.SplittablePayload$WriteRequestEncoder.encode(SplittablePayload.java:221)\n at com.mongodb.internal.connection.SplittablePayload$WriteRequestEncoder.encode(SplittablePayload.java:187)\n at org.bson.codecs.BsonDocumentWrapperCodec.encode(BsonDocumentWrapperCodec.java:63)\n at org.bson.codecs.BsonDocumentWrapperCodec.encode(BsonDocumentWrapperCodec.java:29)\n at com.mongodb.internal.connection.BsonWriterHelper.writeDocument(BsonWriterHelper.java:77)\n at com.mongodb.internal.connection.BsonWriterHelper.writePayload(BsonWriterHelper.java:59)\n at com.mongodb.internal.connection.CommandMessage.encodeMessageBodyWithMetadata(CommandMessage.java:162)\n at com.mongodb.internal.connection.RequestMessage.encode(RequestMessage.java:138)\n at com.mongodb.internal.connection.CommandMessage.encode(CommandMessage.java:59)\n at com.mongodb.internal.connection.InternalStreamConnection.sendAndReceive(InternalStreamConnection.java:268)\n at com.mongodb.internal.connection.UsageTrackingInternalConnection.sendAndReceive(UsageTrackingInternalConnection.java:100)\n at com.mongodb.internal.connection.DefaultConnectionPool$PooledConnection.sendAndReceive(DefaultConnectionPool.java:490)\n at com.mongodb.internal.connection.CommandProtocolImpl.execute(CommandProtocolImpl.java:71)\n at com.mongodb.internal.connection.DefaultServer$DefaultServerProtocolExecutor.execute(DefaultServer.java:253)\n at com.mongodb.internal.connection.DefaultServerConnection.executeProtocol(DefaultServerConnection.java:202)\n at com.mongodb.internal.connection.DefaultServerConnection.command(DefaultServerConnection.java:118)\n at com.mongodb.internal.operation.MixedBulkWriteOperation.executeCommand(MixedBulkWriteOperation.java:431)\n at com.mongodb.internal.operation.MixedBulkWriteOperation.executeBulkWriteBatch(MixedBulkWriteOperation.java:251)\n at com.mongodb.internal.operation.MixedBulkWriteOperation.access$700(MixedBulkWriteOperation.java:76)\n at com.mongodb.internal.operation.MixedBulkWriteOperation$1.call(MixedBulkWriteOperation.java:194)\n at com.mongodb.internal.operation.MixedBulkWriteOperation$1.call(MixedBulkWriteOperation.java:185)\n at com.mongodb.internal.operation.OperationHelper.withReleasableConnection(OperationHelper.java:621)\n at com.mongodb.internal.operation.MixedBulkWriteOperation.execute(MixedBulkWriteOperation.java:185)\n at com.mongodb.internal.operation.MixedBulkWriteOperation.execute(MixedBulkWriteOperation.java:76)\n at com.mongodb.client.internal.MongoClientDelegate$DelegateOperationExecutor.execute(MongoClientDelegate.java:187)\n at com.mongodb.client.internal.MongoCollectionImpl.executeBulkWrite(MongoCollectionImpl.java:442)\n at com.mongodb.client.internal.MongoCollectionImpl.bulkWrite(MongoCollectionImpl.java:422)\n at com.mongodb.kafka.connect.sink.MongoSinkTask.bulkWriteBatch(MongoSinkTask.java:209)\n ... 13 more\n"}],"type":"sink"}
Kafka Topic中的文档内容
{"_id": {"$oid": "634fd99b52281517a468f3a7"},"schema": {"type": "struct", "fields": [{"type": "int32","optional": true, "field": "id"}, {"type": "string", "optional": true, "field": "name"}, {"type": "string", "optional": true, "field": "middel_name"}, {"type": "string", "optional": true, "field": "surname"}],"optional": false, "name": "foobar"},"payload": {"id":45,"name":"mongo","middle_name": "mmp","surname": "kafka"}}
Connector配置
{ "name": "mongodb-sink-connector", "config": { "connector.class": "com.mongodb.kafka.connect.MongoSinkConnector", "topics": "migration-mongo", "connection.uri": "mongodb://abc:xyz@xx.xx.xx.01:27018,xx.xx.xx.02:27018,xx.xx.xx.03:27018/?authSource=admin&replicaSet=dev", "key.converter":"org.apache.kafka.connect.json.JsonConverter", "key.converter.schemas.enable": "false", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter.schemas.enable": "false", "document.id.strategy.overwrite.existing": "false", "validate.non.null": false, "database": "foo", "collection": "product" } }
解决方案
原因分析
当前使用的org.apache.kafka.connect.json.JsonConverter会把$oid当作普通JSON字段处理,但MongoDB的BSON解析器不允许字段名以$开头(除了MongoDB扩展JSON规范中的特殊类型标记),因此触发Invalid BSON field name $oid错误。
方法一:使用MongoDB专属JSON转换器
将连接器配置中的value.converter替换为MongoDB提供的转换器,它能正确识别$oid这类扩展JSON语法,自动转换为BSON的ObjectId类型:
"value.converter": "com.mongodb.kafka.connect.json.MongoJsonConverter", "value.converter.schemas.enable": "false"
修改后的完整配置:
{ "name": "mongodb-sink-connector", "config": { "connector.class": "com.mongodb.kafka.connect.MongoSinkConnector", "topics": "migration-mongo", "connection.uri": "mongodb://abc:xyz@xx.xx.xx.01:27018,xx.xx.xx.02:27018,xx.xx.xx.03:27018/?authSource=admin&replicaSet=dev", "key.converter":"org.apache.kafka.connect.json.JsonConverter", "key.converter.schemas.enable": "false", "value.converter": "com.mongodb.kafka.connect.json.MongoJsonConverter", "value.converter.schemas.enable": "false", "document.id.strategy.overwrite.existing": "false", "validate.non.null": false, "database": "foo", "collection": "product" } }
方法二:修改Kafka文档格式
如果不想更换转换器,可直接将_id字段的值改为字符串格式,MongoDB会自动将其转换为ObjectId类型:
{"_id": "634fd99b52281517a468f3a7","schema": {"type": "struct", "fields": [{"type": "int32","optional": true, "field": "id"}, {"type": "string", "optional": true, "field": "name"}, {"type": "string", "optional": true, "field": "middel_name"}, {"type": "string", "optional": true, "field": "surname"}],"optional": false, "name": "foobar"},"payload": {"id":45,"name":"mongo","middle_name": "mmp","surname": "kafka"}}
推荐方案
如果需要保留_id的ObjectId类型,优先选择方法一,因为MongoJsonConverter专门为MongoDB扩展JSON设计,还能支持$date、$numberLong等其他BSON类型的JSON表示。
内容的提问来源于stack exchange,提问作者Ali Sheer
相关产品推荐
相关产品推荐

