求助:Debezium MongoDB连接器生成记录Payload含反斜杠
Debezium MongoDB连接器
after字段含转义字符问题 当前输出异常现象
使用Debezium MongoDB连接器抽取数据时,整体运行正常,但生成的记录Payload中after属性是带反斜杠的JSON字符串,而非结构化JSON对象,示例如下:
{ "after": "{\"_id\": {\"$oid\": \"63626d5993801d8fd1140993\"},\"document\": \"29973569000204\",\"document_type\": \"CNPJ\"}", "patch": null, "filter": null, "source": { "version": "1.7.1.Final", "connector": "mongodb", "name": "xxxxxxxxxx", "ts_ms": 8466513, "snapshot": "false", "db": "database", "sequence": null, "rs": "atlas-iurhise-shard-0", "collection": "mongo_collection", "ord": 1, "h": null, "tord": 4, "stxnid": "281f4230-d8cc-3d23-a556-89923b45e25f:168" }, "op": "c", "ts_ms": 1667394905422, "transaction": null }
期望输出格式
希望after字段为结构化JSON对象,格式如下:
{ "after": { "_id": { "$oid": "63626d5993801d8fd1140993" }, "document": "29973585214796", "document_type": "CNPJ" }, "patch": null, "filter": null, "source": { "version": "1.7.1.Final", "connector": "mongodb", "name": "xxxxxxxxxx", "ts_ms": 8466513, "snapshot": "false", "db": "database", "sequence": null, "rs": "atlas-iurhise-shard-0", "collection": "mongo_collection", "ord": 1, "h": null, "tord": 4, "stxnid": "281f4230-d8cc-3d23-a556-89923b45e25f:168" }, "op": "c", "ts_ms": 1667394905422, "transaction": null }
初始连接器配置
已尝试相关解决方案但无效,初始配置如下:
{ "name": "DebeziumDataExtract", "config": { "connector.class": "io.debezium.connector.mongodb.MongoDbConnector", "tasks.max": "3", "mongodb.hosts": "removed", "mongodb.name": "removed", "mongodb.user": "removed", "mongodb.password": "removed", "mongodb.ssl.enabled": "true", "collection.whitelist": "removed", "key.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter": "org.apache.kafka.connect.storage.StringConverter", "hstore.handling.mode": "json", "decimal.handling.mode": "string", "key.converter.schemas.enable": "false", "value.converter.schemas.enable": "false", "heartbeat.interval.ms": "1000", "heartbeat.topics.prefix": "removed", "topic.creation.default.replication.factor": 3, "topic.creation.default.partitions": 1, "topic.creation.default.cleanup.policy": "compact", "topic.creation.default.compression.type": "lz4", "transforms": "unwrap", "transforms.unwrap.collection.expand.json.payload": "true" } }
更新后的配置
根据建议调整配置后问题仍未解决,更新后的配置如下:
{ "name": "DebeziumTransportPlanner", "config": { "connector.class": "io.debezium.connector.mongodb.MongoDbConnector", "tasks.max": "3", "mongodb.hosts": "stg-transport-planner-0-shard-00-00-00.xmapa.mongodb.net,stg-transport-planner-0-shard-00-01.xmapa.mongodb.net,stg-transport-planner-0-shard-00-02.xmapa.mongodb.net", "mongodb.name": "stg-transport-planner-01", "mongodb.user": "oploguser-stg", "mongodb.password": "vCh1NtV4PoY8PeSJ", "mongodb.ssl.enabled": "true", "collection.whitelist": "stg-transport-planner-01[.]aggregated_transfers", "key.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "hstore.handling.mode": "json", "decimal.handling.mode": "string", "key.converter.schemas.enable": "false", "value.converter.schemas.enable": "false", "heartbeat.interval.ms": "1000", "heartbeat.topics.prefix": "__debeziumtransport-planner-heartbeat", "topic.creation.default.replication.factor": 3, "topic.creation.default.partitions": 1, "topic.creation.default.cleanup.policy": "compact", "topic.creation.default.compression.type": "lz4", "transforms": "unwrap", "transforms.unwrap.type":"io.debezium.connector.mongodb.transforms.ExtractNewDocumentState", "transforms.unwrap.collection.expand.json.payload": "true", "transforms.unwrap.collection.fields.additional.placement": "route_external_id:header,transfer_index:header" } }
内容的提问来源于stack exchange,提问作者João Zarate
相关产品推荐
相关产品推荐

