如何让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
相关产品推荐
相关产品推荐

