如何在Apache NiFi中获取输入队列的FlowFile数量并作为参数读取
在Apache NiFi中获取输入队列FlowFile数量并传递为参数
一、查看输入队列的FlowFile数量
1. NiFi UI直接查看
每个处理器之间的连接(即输入队列)会在UI上实时显示当前的FlowFile数量,鼠标悬停还能查看更详细的统计信息(如总大小、排队时长等)。
2. 程序化获取(API方式)
如果需要通过代码或流程内部获取,可以调用NiFi的REST API:
- 请求地址:
GET /nifi-api/flow/connections/{connection-id} - 返回的JSON响应中,
aggregateSnapshot.flowFileCount字段即为队列的FlowFile总数。
二、将队列数量作为参数传递给后续处理器
要在流程中把指定队列的FlowFile数量提取为参数,可通过以下两种方式实现:
方式一:REST API + 解析处理器组合
步骤1:获取目标队列的ID
在NiFi UI中点击目标连接(输入队列),右侧面板的「Settings」标签下可以找到该队列的唯一ID,复制备用。
步骤2:配置InvokeHTTP调用API
- HTTP Method:选择
GET - Remote URL:填入你的NiFi地址+队列ID,例如
http://localhost:8080/nifi-api/flow/connections/abc123-xxx-yyy - 若NiFi开启了认证,需在「Authentication」标签下配置对应认证方式(如Basic Auth)。
步骤3:用EvaluateJsonPath解析结果
- Destination:选择
FlowFile Attribute - 添加自定义属性,例如命名为
queue-flowfile-count,属性值设置为$.aggregateSnapshot.flowFileCount - 处理器执行后,FlowFile会新增
queue-flowfile-count属性,值为队列的FlowFile数量。
步骤4:后续处理器引用参数
在后续任意处理器中,通过NiFi表达式语言${queue-flowfile-count}即可直接读取该数值作为参数使用。
方式二:ExecuteScript直接调用内部API(更高效)
通过Groovy脚本直接调用NiFi内部服务,无需走外部REST API,效率更高:
- 添加
ExecuteScript处理器,选择Groovy作为脚本语言 - 替换脚本中的
your-connection-id为目标队列的ID,脚本内容如下:
import org.apache.nifi.controller.FlowController import org.apache.nifi.controller.Connection // 获取FlowController实例 def flowController = context.controllerServiceLookup.getControllerServices(FlowController).find().get() // 根据ID获取目标连接(队列) def targetConnection = flowController.flowManager.rootGroup.findConnectionById("your-connection-id") // 获取队列中的FlowFile数量 def flowFileCount = targetConnection.flowFileQueue.activeQueueSize.objectCount.toString() // 将数量存入FlowFile属性 flowFile = session.putAttribute(flowFile, "queue-flowfile-count", flowFileCount) REL_SUCCESS << flowFile
- 后续处理器同样通过
${queue-flowfile-count}引用该参数。
注意事项
- 若为NiFi集群部署,REST API地址需指向集群节点或负载均衡入口
- 调用API或内部服务的账号需具备访问队列信息的权限
- 如需定时获取队列数量,可搭配
GenerateFlowFile处理器定时触发流程
内容的提问来源于stack exchange,提问作者Amarnatha Reddy
相关产品推荐
相关产品推荐

