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

如何配置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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 12:03:23