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

如何不使用HeaderToField SMT将Kafka Headers中的Auditlogs保存到MongoDB Sink?

如何在不使用HeaderToField SMT的情况下将Kafka Header中的Auditlogs值存入MongoDB Sink
  • 修改生产者逻辑,提前把Header值塞进消息体
    如果能改动消息生产者的代码,发送消息时直接把Auditlogs Header的值取出来,作为消息体的一个字段一起发送。这样MongoDB Sink消费消息时,直接从消息体里读这个字段就能存进数据库,完全不用依赖任何SMT。

  • 给MongoDB Sink写自定义Converter
    MongoDB Sink支持自定义Converter实现,你可以自己写一个Converter,在处理消息的时候主动读取Kafka Header里的Auditlogs值,把它加到要写入MongoDB的文档里:

    • 实现org.apache.kafka.connect.storage.Converter接口,重点重写toConnectData方法,在这个方法里拿到Kafka记录的Headers,提取Auditlogs的值。
    • 把提取到的值合并到Connect的数据结构(比如Struct或者Map)里。
    • 在MongoDB Sink的配置里指定这个自定义Converter的全类名,替换掉默认的Converter就行。
  • 用Kafka Streams先预处理消息
    在MongoDB Sink前面加个Kafka Streams应用当中间层:

    • 读取原始Kafka主题的消息,提取Auditlogs Header的值。
    • 把这个值加到消息体里,再把处理后的消息发到一个新的Kafka主题。
    • 让MongoDB Sink消费这个新主题的消息,直接从消息体里拿Auditlogs值存进MongoDB。

内容的提问来源于stack exchange,提问作者sai jyothsna pentyala

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 21:42:41