寻求可将Kafka Header注入消息体的Kafka Connect SMT(AWS MSK适用)
在AWS MSK Connect中提取Kafka Headers字段注入消息体
你可以直接使用Kafka Connect内置的InsertField转换器,通过配置从Header提取字段的参数实现需求,无需自定义转换器。
核心配置示例
替换原有静态字段配置,改用header.field和target.field参数指定要读取的Header键名及注入到消息体的字段名:
"transforms": "InsertHeaderField", "transforms.InsertHeaderField.type": "org.apache.kafka.connect.transforms.InsertField$Value", "transforms.InsertHeaderField.header.field": "idempotence-key", "transforms.InsertHeaderField.target.field": "idempotence-key"
参数说明
transforms.InsertHeaderField.header.field:指定要读取的Kafka Header键名(比如示例中的idempotence-key)transforms.InsertHeaderField.target.field:指定注入到消息体的字段名(可与Header键名一致,也可自定义)
多Header字段注入
如果需要同时注入多个Header字段,添加多个转换器实例即可:
"transforms": "InsertIdempotenceKey,InsertSource", "transforms.InsertIdempotenceKey.type": "org.apache.kafka.connect.transforms.InsertField$Value", "transforms.InsertIdempotenceKey.header.field": "idempotence-key", "transforms.InsertIdempotenceKey.target.field": "idempotence-key", "transforms.InsertSource.type": "org.apache.kafka.connect.transforms.InsertField$Value", "transforms.InsertSource.header.field": "source", "transforms.InsertSource.target.field": "source"
效果验证
按示例配置后,原消息体和Header会被转换为:
{ "id": 24, "ts": 1626102708861, "name": "John Smith", "book": "Kafka: The Definitive Guide", "idempotence-key":"019a6d9f-9538t" }
该方案完全兼容AWS MSK Connect,无需额外依赖,直接使用内置转换器即可实现需求。
内容的提问来源于stack exchange,提问作者Sylvester
相关产品推荐
相关产品推荐

