Debezium MongoDB源连接器为何生成string类型after字段而非JSON对象?
Debezium MongoDB连接器payload.after为字符串导致unwrap转换失败问题
问题背景
使用Debezium MongoDB源连接器(版本3.0.6.Final),配置如下:
{ "name": "mongo-debezium-connector", "config": { "connector.class": "io.debezium.connector.mongodb.MongoDbConnector", "tasks.max": "1", "mongodb.connection.string": "mongodb://mongo:27017/?replicaSet=rs0", "database.include.list": "sample", "collection.include.list": "sample.workflows,sample.simulations", "topic.prefix": "mongo", "key.converter": "org.apache.kafka.connect.storage.StringConverter", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter.schemas.enable": true } }
插入MongoDB的示例文档:
{ "_id": { "$oid": "676d8e51105e01702fe9496c" }, "name": "workflow 3" }
通过Kowl查看Kafka事件时,发现payload.after字段为字符串类型而非JSON对象:
{ "schema": { "type": "struct", "fields": [ { "type": "string", "optional": true, "name": "io.debezium.data.Json", "version": 1, "field": "before" }, { "type": "string", "optional": true, "name": "io.debezium.data.Json", "version": 1, "field": "after" }, { "type": "struct", "fields": [ { "type": "array", "items": { "type": "string", "optional": false }, "optional": true, "field": "removedFields" }, { "type": "string", "optional": true, "name": "io.debezium.data.Json", "version": 1, "field": "updatedFields" }, { "type": "array", "items": { "type": "struct", "fields": [ { "type": "string", "optional": false, "field": "field" }, { "type": "int32", "optional": false, "field": "size" } ], "optional": false, "name": "io.debezium.connector.mongodb.changestream.truncatedarray", "version": 1 }, "optional": true, "field": "truncatedArrays" } ], "optional": true, "name": "io.debezium.connector.mongodb.changestream.updatedescription", "version": 1, "field": "updateDescription" }, { "type": "struct", "fields": [ { "type": "string", "optional": false, "field": "version" }, { "type": "string", "optional": false, "field": "connector" }, { "type": "string", "optional": false, "field": "name" }, { "type": "int64", "optional": false, "field": "ts_ms" }, { "type": "string", "optional": true, "name": "io.debezium.data.Enum", "version": 1, "parameters": { "allowed": "true,first,first_in_data_collection,last_in_data_collection,last,false,incremental" }, "default": "false", "field": "snapshot" }, { "type": "string", "optional": false, "field": "db" }, { "type": "string", "optional": true, "field": "sequence" }, { "type": "int64", "optional": true, "field": "ts_us" }, { "type": "int64", "optional": true, "field": "ts_ns" }, { "type": "string", "optional": false, "field": "collection" }, { "type": "int32", "optional": false, "field": "ord" }, { "type": "string", "optional": true, "field": "lsid" }, { "type": "int64", "optional": true, "field": "txnNumber" }, { "type": "int64", "optional": true, "field": "wallTime" } ], "optional": false, "name": "io.debezium.connector.mongo.Source", "field": "source" }, { "type": "string", "optional": true, "field": "op" }, { "type": "int64", "optional": true, "field": "ts_ms" }, { "type": "struct", "fields": [ { "type": "string", "optional": false, "field": "id" }, { "type": "int64", "optional": false, "field": "total_order" }, { "type": "int64", "optional": false, "field": "data_collection_order" } ], "optional": true, "name": "event.block", "version": 1, "field": "transaction" } ], "optional": false, "name": "mongo.sample.workflows.Envelope" }, "payload": { "before": null, "after": "{\"_id\": {\"$oid\": \"676d8e51105e01702fe9496c\"},\"name\": \"workflow 3\"}", "updateDescription": null, "source": { "version": "3.0.6.Final", "connector": "mongodb", "name": "mongo", "ts_ms": 1735233105000, "snapshot": "false", "db": "sample", "sequence": null, "ts_us": 1735233105000000, "ts_ns": 1735233105000000000, "collection": "workflows", "ord": 1, "lsid": null, "txnNumber": null, "wallTime": 1735233105571 }, "op": "c", "ts_ms": 1735233105616, "transaction": null } }
尝试应用ExtractNewRecordState类型的unwrap转换时,抛出以下错误:
org.apache.kafka.connect.errors.DataException: Only Struct objects supported for [source field insertion], found: java.lang.String
问题原因及排查方向
核心原因
Debezium MongoDB连接器默认将MongoDB文档序列化为JSON字符串存入before/after字段,而非结构化的Struct类型。这是因为MongoDB是无Schema数据库,连接器无法提前知晓文档结构,因此默认用JSON字符串保留原始数据格式。而ExtractNewRecordState转换要求处理结构化的Struct对象,因此报错。
排查方向
- 启用Schema推断:在连接器配置中添加
mongodb.schema.inference.enabled=true,让连接器自动推断文档结构,将before/after字段生成为Struct类型。注意集合中文档结构不一致时,推断出的Schema会包含可选字段,需确保下游能兼容。 - 检查转换器配置:当前使用的
JsonConverter启用了Schema,但仅当连接器输出Struct类型时,转换器才会生成结构化数据。不要关闭value.converter.schemas.enable,否则会丢失Schema信息,仍无法被unwrap转换处理。 - 验证集合结构一致性:若集合中文档结构差异极大,Schema推断可能无法生成稳定的Struct类型,导致部分记录仍为字符串。此时需评估是否统一文档结构,或使用自定义转换逻辑处理字符串类型的
after字段。 - 检查版本兼容性:确认3.0.6.Final版本是否存在Schema推断相关Bug,必要时升级到最新稳定版。
内容的提问来源于stack exchange,提问作者Rasheduzzaman Sourov
相关产品推荐
相关产品推荐

