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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 08:57:02