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

如何在MongoDB Sink的文档中保存Kafka消息Key?

问题

使用MongoDB Sink Connector时,已成功将Kafka AVRO消息的Value保存到MongoDB文档,但需要同时将Kafka消息的Key也存入文档。尝试用org.apache.kafka.connect.transforms.HoistField$Key将Key添加到Value中未生效;使用ProvidedInKeyStrategy虽能将Key写入文档,但会把Key设为MongoDB文档的_id,不符合需求。

解决方案

通过两个Kafka Connect Transform的组合实现:先将Kafka Key包装为独立字段,再合并到Value中,同时保留默认的_id生成逻辑。

修改后的完整配置如下:

"config": {
    "connector.class": "com.mongodb.kafka.connect.MongoSinkConnector",
    "connection.uri": "mongodb://mongo1",
    "database": "mongodb",
    "collection": "sink",
    "topics": "topics.foo",

    "value.converter": "io.confluent.connect.avro.AvroConverter",
    "value.converter.schema.registry.url": "http://schema-registry:8081",
    "key.converter": "io.confluent.connect.avro.AvroConverter",
    "key.converter.schema.registry.url": "http://schema-registry:8081",
    
    "transforms": "hoistKey,mergeKeyIntoValue",

    // 将Kafka Key包装为名为kafkaKey的字段
    "transforms.hoistKey.type": "org.apache.kafka.connect.transforms.HoistField$Key",
    "transforms.hoistKey.field": "kafkaKey",

    // 将Key中的kafkaKey字段合并到Value中
    "transforms.mergeKeyIntoValue.type": "org.apache.kafka.connect.transforms.Merge$Value",
    "transforms.mergeKeyIntoValue.fields": "kafkaKey",
    "transforms.mergeKeyIntoValue.delete.fields": "kafkaKey"
  }
配置说明
  • HoistField$Key:把原始的Kafka Key(复杂AVRO结构)包装成kafkaKey字段,此时Key的结构变为{"kafkaKey": 原始Key内容}。
  • Merge$Value:将Key中kafkaKey字段的内容合并到Value里,同时删除Key中的该字段(可选,避免冗余)。最终Value会包含原有timestamp字段和新增的kafkaKey字段,MongoDB文档完整保存这些内容,_id保持默认自动生成逻辑。
验证结果

保存后的MongoDB文档结构示例:

{
  "_id": ObjectId("..."),
  "timestamp": 1620000000000,
  "kafkaKey": {
    "conversation_id": "xxx",
    "broker_key": "yyy",
    "user_id": "zzz",
    "application": "app1",
    "environment": "Dev"
  }
}

内容的提问来源于stack exchange,提问作者Zikria Azimi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 18:20:47