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

寻求可将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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 01:39:58