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

如何通过Kafka Connect提取消息Key并追加至Value中

解决Kafka Connect将消息Key追加到Value的SMT配置问题

你当前的配置错误在于使用了InsertField$Value的静态字段插入参数,而非提取消息Key的参数。要实现把消息Key追加到Value中的需求,只需修改SMT的配置项即可:

修正后的连接器配置

{
    "name": "aaa",
    "config": {
        "name": "aaa",
        "connector.class": "io.aiven.kafka.connect.opensearch.OpensearchSinkConnector",
        "tasks.max": "1",
        "topics": "test_x",
        "transforms": "InsertKeyToValue",
        "transforms.InsertKeyToValue.type": "org.apache.kafka.connect.transforms.InsertField$Value",
        "transforms.InsertKeyToValue.key.field": "key",
        "key.ignore": "false",
        "value.converter": "org.apache.kafka.connect.json.JsonConverter",
        "value.converter.schemas.enable": "false",
        "key.converter.schemas.enable": "false",
        "key.converter": "org.apache.kafka.connect.storage.StringConverter",
        "connection.url": "https://localhost:9200",
        "connection.username": "test",
        "connection.password": "test",
        "batch.size": "10000",
        "linger.ms": "1000",
        "errors.tolerance": "all",
        "errors.log.include.messages": "true",
        "errors.log.enable": "true"
    }
}

关键配置说明

  • 移除原配置中static.field和static.value这两个插入固定值的参数,替换为key.field": "key":该参数会指示InsertField$Value SMT提取消息Key的值,插入到Value的key字段中。
  • 保留key.converter": "org.apache.kafka.connect.storage.StringConverter":你的消息Key是字符串类型,该转换器可正确解析Key的值。
  • key.ignore": "false"确保连接器不会忽略消息Key,保证SMT能获取到Key的值。

配置生效后,发送的消息Value {"name":"ali"}和Key Person会被处理成{"name":"ali","key":"Person"}并写入OpenSearch。

内容的提问来源于stack exchange,提问作者Alihossein shahabi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 16:40:02