如何不使用HeaderToField SMT将Kafka Headers中的Auditlogs保存到MongoDB Sink?
如何在不使用HeaderToField SMT的情况下将Kafka Header中的Auditlogs值存入MongoDB Sink
修改生产者逻辑,提前把Header值塞进消息体
如果能改动消息生产者的代码,发送消息时直接把AuditlogsHeader的值取出来,作为消息体的一个字段一起发送。这样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主题的消息,提取
AuditlogsHeader的值。 - 把这个值加到消息体里,再把处理后的消息发到一个新的Kafka主题。
- 让MongoDB Sink消费这个新主题的消息,直接从消息体里拿
Auditlogs值存进MongoDB。
- 读取原始Kafka主题的消息,提取
内容的提问来源于stack exchange,提问作者sai jyothsna pentyala
相关产品推荐
相关产品推荐

