如何在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
相关产品推荐
相关产品推荐

