能否通过Kafka Connect将Kafka包装JSON消息拆分后按类型存入S3?
我通过Kafka接收的是JSON格式的包装负载,示例如下:
{"format": "wrapper","time": 1626814608000,"events": [{"id": "item1","type": "product1","count": 200},{"id": "item2","type": "product2","count": 300}],"metadata": {"schema": "schema-1"}}
我需要把这些数据导出到S3,但不能保留包装结构,要将每个事件按其type字段分别存储到对应前缀下。比如:
bucket/product1路径下存储:{"id": "item1", "type": "product1", "count": 200}bucket/product2路径下存储:{"id": "item2", "type": "product2", "count": 300}
我的疑问是:能否用Kafka Connect实现这个需求?
我了解Kafka Connect的Single Message Transforms(SMT)只能做单消息转换(签名为R=>R),不支持把单条消息拆分成多条的扇出操作。目前我觉得用Kafka Connect直接实现不了,但想确认有没有遗漏的方法,再考虑其他方案。
你说的没错,原生Kafka Connect的SMT确实不支持单消息拆分成多条的扇出操作,但可以通过以下几种方式实现需求:
使用第三方拆分插件
社区有现成的Kafka Connect插件可以实现消息拆分,比如kafka-connect-splitter,它能将包含数组的包装消息拆分成多条独立消息,每条对应数组里的一个事件元素。拆分完成后,再用SMT提取事件的type字段,配置S3 Sink Connector时将该字段设为存储前缀即可。自定义Transform插件
如果你具备开发能力,可以自行实现一个支持扇出的Transform插件。虽然官方的Transformation接口定义是单输入单输出,但可以通过扩展逻辑返回多条SourceRecord来实现消息拆分——不过要注意兼容Kafka Connect的运行机制,保证任务的容错性和数据一致性。结合Kafka Streams预处理
先编写一个轻量的Kafka Streams应用,消费原始主题的包装消息,拆分出每个事件后发送到临时主题(也可以按type字段直接分主题),再用Kafka Connect的S3 Sink从这些临时主题读取数据,配置前缀规则按type字段存储到S3。这种方式无需修改Connect插件,实现起来更灵活。
内容的提问来源于stack exchange,提问作者ferrari16

