NiFi 2.2.0中如何为Notify&Wait流程配置缓存服务?
NiFi 2.x 配置 Notify & Wait 实现未知数量FlowFile批量完成监控
NiFi 2.x 已移除旧版的DistributedMapCacheServer,改用MapCacheServer及对应客户端服务实现缓存协同。以下是适配新版本的Notify&Wait配置步骤,实现未知数量FlowFile全部处理完成的监控:
1. 准备缓存服务组件
- 服务端:创建
MapCacheServer的具体实现(如HazelcastMapCacheServer支持分布式集群,CaffeineMapCacheServer为本地缓存),设置自定义缓存名称(如batch-monitor-cache),按需调整其他参数(如集群地址、缓存大小)。 - 客户端:创建对应类型的
MapCacheClientService(如HazelcastMapCacheClientService),关联上述MapCacheServer,确保缓存名称与服务端完全一致,分布式场景下需配置正确的服务端地址。
2. 配置Notify处理器(区分数据与终止信号)
由于FlowFile数量未知,需通过终止信号标记批次完成:
- 先通过
RouteOnAttribute等组件,将业务数据FlowFile和手动生成的终止信号FlowFile(可通过GenerateFlowFile生成)分开路由:- 针对数据FlowFile的Notify配置:
- 选择
Cache Service为已配置的MapCacheClientService Release Signal Identifier设置为统一批次标识(如${batch.id},需确保所有同批次FlowFile带有该属性)Signal Action选择Increment Cache Value,Signal Count设为1(每处理一个数据FlowFile,缓存计数+1)
- 选择
- 针对终止信号FlowFile的Notify配置:
- 选择同一个
MapCacheClientService Release Signal Identifier与数据FlowFile的配置一致Signal Action选择Mark Signal For Termination(发送批次完成的终止标记)
- 选择同一个
- 针对数据FlowFile的Notify配置:
3. 配置Wait处理器
- 选择
Cache Service为同一个MapCacheClientService Release Signal Identifier与Notify的批次标识完全一致Wait Strategy选择Wait For Termination Signal- 设置
Maximum Wait Time(如3600秒)避免无限等待 - 当所有数据FlowFile完成Notify计数,且终止信号的Notify发送终止标记后,Wait处理器会释放对应的FlowFile,确认批次全部处理完成
常见问题排查
- 若Wait无法识别连接:
- 确认Notify与Wait使用完全相同的MapCacheClientService,且服务已启用(状态为有效)
- 检查
Release Signal Identifier的表达式是否一致,且FlowFile确实包含该属性 - 验证MapCacheServer与Client的配置匹配(缓存名称、地址等参数一致)
内容的提问来源于stack exchange,提问作者Mike2116 -PRO
相关产品推荐
相关产品推荐

