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

如何在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,效率更高:

  1. 添加ExecuteScript处理器,选择Groovy作为脚本语言
  2. 替换脚本中的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
  1. 后续处理器同样通过${queue-flowfile-count}引用该参数。

注意事项

  • 若为NiFi集群部署,REST API地址需指向集群节点或负载均衡入口
  • 调用API或内部服务的账号需具备访问队列信息的权限
  • 如需定时获取队列数量,可搭配GenerateFlowFile处理器定时触发流程

内容的提问来源于stack exchange,提问作者Amarnatha Reddy

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 03:05:00