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

如何让Kafka SQS Source Connector输出纯JSON格式消息?

解决Kafka SQS Source Connector输出纯JSON消息的问题

方法1:使用Single Message Transform(SMT)提取Body字段

在连接器配置中添加以下SMT规则,直接提取Struct中的Body字段作为Kafka消息的value:

transforms=ExtractBody
transforms.ExtractBody.type=org.apache.kafka.connect.transforms.ExtractField$Value
transforms.ExtractBody.field=Body

配置生效后,连接器会自动丢弃Struct中的元数据字段(如ApproximateFirstReceiveTimestamp、SenderId等),仅保留Body内的内容作为消息值。

方法2:配置JSON转换器并启用Body解析

若SQS消息Body本身为JSON格式,可通过以下配置让连接器直接输出纯JSON字符串:

value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=false
sqs.body.json.parser.enabled=true
  • value.converter指定为JSON转换器,避免消息被序列化为Struct格式
  • value.converter.schemas.enable=false禁用Schema输出,确保消息为纯JSON字符串
  • sqs.body.json.parser.enabled=true让连接器自动解析SQS Body中的JSON内容,直接输出目标结构

针对JSON数组Body的补充处理

如果SQS Body是JSON数组(如示例中的[{...}]),需要提取数组内的单个对象时,可追加额外SMT配置:

transforms=ExtractBody,ExtractArrayItem
transforms.ExtractBody.type=org.apache.kafka.connect.transforms.ExtractField$Value
transforms.ExtractBody.field=Body
transforms.ExtractArrayItem.type=org.apache.kafka.connect.transforms.ExtractField$Value
transforms.ExtractArrayItem.field=0

此配置会先提取Body数组,再取出数组的第一个元素作为最终消息内容。

内容的提问来源于stack exchange,提问作者Mike Marsh

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 01:45:45