如何配置Kafka MongoDB源连接器以发布原始JSON而非字符串
问题
当前使用Kafka MongoDB源连接器,连接器配置如下:
{ "output.format.key":"json", "output.format.value":"json", "copy.existing":"false", "publish.full.document.only":"true", "change.stream.full.document":"updateLookup", "output.schema.infer.value":"true", "value.converter":"org.apache.kafka.connect.json.JsonConverter", "value.converter.schemas.enable":"false", "key.converter":"org.apache.kafka.connect.json.JsonConverter", "key.converter.schemas.enable":"false" }
connect-distributed.properties全局配置如下:
key.converter=org.apache.kafka.connect.json.JsonConverter value.converter=org.apache.kafka.connect.json.JsonConverter key.converter.schemas.enable=false value.converter.schemas.enable=false
目前Kafka主题中获取到的值是字符串包裹的JSON格式,示例:
"{\"_id\": {\"$oid\": \"67f59ce3f73ca222999beab5\"}, \"product\": \"Iphone 213\", \"price\": 9009}"
需要输出原始JSON格式,而非字符串包裹的形式。
解决方案
问题根源是连接器配置与全局配置重复设置了JSON转换器,导致MongoDB连接器输出的JSON被JsonConverter再次序列化,最终变成字符串格式。
修改步骤:
- 移除连接器配置中的
value.converter、value.converter.schemas.enable、key.converter、key.converter.schemas.enable这四个配置项,全局配置已经定义了这些规则,无需在连接器中重复配置。 - 保留
output.format.value=json,这是MongoDB连接器用来生成JSON格式数据的核心配置,配合全局的JsonConverter(关闭schema)就能输出原始JSON。
修改后的连接器配置如下:
{ "output.format.key":"json", "output.format.value":"json", "copy.existing":"false", "publish.full.document.only":"true", "change.stream.full.document":"updateLookup", "output.schema.infer.value":"true" }
修改完成后,Kafka主题中的消息值会变成原始JSON格式:
{"_id": {"$oid": "67f59ce3f73ca222999beab5"}, "product": "Iphone 213", "price": 9009}
内容的提问来源于stack exchange,提问作者RRR
相关产品推荐
相关产品推荐

