如何在Apache NiFi中从ConsumeKafka提取键值对,去重后通过PutKafka发送至其他Topic?
我来帮你拆解这个NiFi+Kafka的需求,从消费、提取键值对、去重到转发的全流程都给你梳理清楚,都是实际生产中验证过的方案:
整体流程概览
整个流程的组件链路大概是:ConsumeKafkaRecord_2_6(或对应版本的ConsumeKafka) → 提取键值对组件(如EvaluateJsonPath/ExtractText) → 去重组件(DistinctHashLookup+缓存服务) → PutKafkaRecord_2_6(或对应版本的PutKafka)
1. 配置ConsumeKafka消费源Topic
首先搞定消息消费的基础配置:
- 核心配置:
- 设置
Bootstrap Servers为你的Kafka集群地址 Topic Name填你要消费的源TopicConsumer Group ID设一个唯一的ID,确保消费组的偏移量管理正常
- 设置
- 偏移量策略:如果要配合去重逻辑,建议选
Manual Offset Management,这样只有当消息成功处理并完成去重校验后,再提交偏移量,避免因故障导致的重复消费。当然如果用Kafka自带的偏移量存储(Kafka Offset Storage)也可以,但手动管理更稳妥。
2. 提取键值对
这一步取决于你的消息格式,分两种常见情况:
情况1:消息是JSON格式
用EvaluateJsonPath组件:
- 添加两个属性,比如:
- 键的提取:
kafka.key=$.your_key_field(把JSON里的目标键字段值存到kafka.key属性) - 值的提取:
kafka.value=$.your_value_field(同理,把值存到kafka.value属性)
- 键的提取:
- 设置
Destination为flowfile-attribute,这样提取的内容会作为FlowFile的属性,方便后续组件调用。
情况2:消息是纯文本/CSV等格式
用ExtractText组件,通过正则表达式提取键值对:
- 比如你的消息格式是
key:xxx,value:yyy,就可以设置:kafka.key=key:(.*?),kafka.value=value:(.*)
- 同样把提取结果存到FlowFile属性中。
3. 基于特定键实现去重(核心步骤)
NiFi里实现分布式/单机去重最可靠的方式是缓存服务+去重组件:
第一步:配置缓存服务
根据你的部署场景选择:
- 单机测试/小流量:用
H2CacheService,设置Cache Directory为NiFi可读写的路径,开启持久化,避免重启后缓存丢失。 - 集群生产环境:用
RedisCacheService,配置Redis集群地址、密码等,天然支持分布式缓存,解决多NiFi节点间的去重一致性问题。
第二步:去重逻辑实现
用DistinctHashLookup组件:
- 配置
Cache Service为你刚才创建的缓存服务 - 设置
Hash Value为${kafka.key}(用之前提取的键作为去重标识) - 组件会自动检查当前键是否在缓存中:
- 如果不存在:路由到
Unmatched关系,同时自动把该键存入缓存,后续消息可以继续流转到PutKafka - 如果存在:路由到
Matched关系,你可以把这些重复消息路由到LogAttribute记录日志,或者直接Drop丢弃
- 如果不存在:路由到
额外优化:缓存过期策略
如果你的键不需要永久去重(比如只需要去重最近24小时的消息),可以在缓存服务里设置Expiration时间,比如86400 seconds,避免缓存无限膨胀占用资源。
4. 用PutKafka发送到目标Topic
最后配置转发组件:
- 设置
Bootstrap Servers为目标Kafka集群地址 Topic Name填目标Topic- 关键配置:
Kafka Key选择Attribute Value,值填${kafka.key}Kafka Value选择Attribute Value,值填${kafka.value}
- 如果需要保证消息顺序,可以开启
Producer Properties里的acks=all,确保消息被Kafka集群确认后再返回成功。
一些踩坑提醒
- 偏移量提交:如果用手动偏移量管理,记得在消息成功发送到目标Topic后,用
CommitKafkaOffset组件提交偏移量,避免重复消费。 - 缓存一致性:集群环境下绝对不要用单机缓存服务(比如H2),否则不同NiFi节点的缓存不共享,会导致重复消息漏判。
- 性能调优:如果消息量很大,可以调整
ConsumeKafka的Max Poll Records和组件的Concurrent Tasks,同时确保缓存服务的性能足够(比如Redis用集群模式)。
内容的提问来源于stack exchange,提问作者hasherBaba
相关产品推荐
相关产品推荐

