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

如何在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填你要消费的源Topic
    • Consumer 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 06:30:03