NiFi按指定属性分组数据遇队列统计与过滤问题求助
背景
需每小时按指定属性对数据分组,单轮队列规模30-40k,极端场景可达200k。MergeContent因无最小/最大数量限制不适用,RouteOnAttribute因属性组合过多不适用。
尝试过的方案及问题
方案1:批量消费分组
尝试消费所有FlowFile并按属性分组生成新FlowFile,但session.getQueueSize().getObjectCount()始终返回10k,即使已调高输出流队列阈值。
方案2:单FlowFile触发过滤匹配
通过自定义逻辑过滤匹配指定属性的FlowFile,代码如下:
final List<FlowFile> flowFiles = session.get(file -> { if (correlationId.equals(Arrays.stream(keys).map(file::getAttribute).collect(Collectors.joining(":")))) return FlowFileFilter.FlowFileFilterResult.ACCEPT_AND_CONTINUE; return FlowFileFilter.FlowFileFilterResult.REJECT_AND_CONTINUE; });
队列含33k个FlowFile时,预期生成约200个分组FlowFile,但实际生成320个,疑似未扫描全部队列。
核心问题
- 是否有参数可调整使
getObjectCount()支持最高300k的队列统计? - 是否可通过调整参数或更换处理器实现全队列FlowFile过滤?
已尝试修改nifi.properties中默认队列阈值为300k,但无效果。
解决方案
针对问题1:调整队列统计上限
session.getQueueSize().getObjectCount()返回10k是因为NiFi默认限制了队列统计的采样上限,由nifi.queue.size.reporting.max.attribute.count参数控制,默认值为10000。需修改nifi.properties:
nifi.queue.size.reporting.max.attribute.count=300000
修改后重启NiFi即可获取准确的大队列数量。注意该参数会影响队列统计性能,过大可能增加节点负载,需根据实际场景调整。
针对问题2:实现全队列FlowFile过滤
方案2未扫描全队列,是因为session.get(FlowFileFilter)受处理器**"Max Flow Files Per Run"**参数限制,默认值通常为10000,导致每次仅扫描部分队列就停止。可通过以下方式解决:
调整处理器运行参数
在自定义处理器配置页,将**"Max Flow Files Per Run"(单轮最大处理FlowFile数)设为大于队列最大规模的值(如300000),同时调整"Max Content Size Per Run"**避免因内容过大触发限制。改用
ListFlowFile+脚本组合
先用ListFlowFile处理器列出队列所有FlowFile的属性(需将Batch Size设为足够大的值),再通过ExecuteScript/ExecuteGroovyScript按属性分组,最后用FetchFlowFile获取对应FlowFile合并。该方式更适配大规模队列分组场景,规避自定义处理器的session限制。优化自定义处理器逻辑
若坚持使用自定义处理器,需在onTrigger方法中循环调用session.get(),直到无匹配FlowFile:
List<FlowFile> allMatchingFlowFiles = new ArrayList<>(); List<FlowFile> batch; do { batch = session.get(file -> { if (correlationId.equals(Arrays.stream(keys).map(file::getAttribute).collect(Collectors.joining(":")))) { return FlowFileFilter.FlowFileFilterResult.ACCEPT_AND_CONTINUE; } return FlowFileFilter.FlowFileFilterResult.REJECT_AND_CONTINUE; }); allMatchingFlowFiles.addAll(batch); } while (!batch.isEmpty()); // 后续处理allMatchingFlowFiles
同时需确保处理器的Max Flow Files Per Run设为足够大的值,避免循环提前终止。
内容的提问来源于stack exchange,提问作者msdalp

