NiFi 1.11:如何让Groovy脚本实时生成新FlowFile会话?
解决NiFi ExecuteScript生成FlowFile延迟入队的问题
核心原因
NiFi的session操作默认会在脚本执行完毕后统一提交,所以你当前生成的5万个FlowFile会全部攒到脚本结束才一次性进入下游队列。
解决方案:分批提交Session
修改Groovy脚本,每生成一定数量的FlowFile就手动提交一次session,这样每批生成的FlowFile会立即进入队列,无需等待整个脚本执行完成。示例代码如下:
// 定义每批提交的数量,可根据实际情况调整 def batchSize = 1000 def count = 0 // 替换为你实际的XML解析/生成循环逻辑 yourXmlRecordList.each { recordOut -> def newFlowFile = session.create() newFlowFile = session.write(newFlowFile, { inStream, outStream -> outStream.write(recordOut.getBytes(StandardCharsets.UTF_8)) } as StreamCallback) session.transfer(newFlowFile, REL_SUCCESS) count++ // 达到批次阈值时提交session if (count % batchSize == 0) { session.commit() // 提交后创建新session,避免会话状态异常 session = sessionFactory.createSession() } } // 提交最后一批不足batchSize的FlowFile if (count % batchSize != 0) { session.commit() }
注意事项
- 批次大小需合理调整:太小会增加NiFi的提交开销,太大则仍会有明显延迟,建议从1000-5000的范围测试适配。
- 异常处理:在提交session时建议用
try-catch包裹,提交失败时调用session.rollback()避免数据异常。 - 不要保留旧session引用:提交后必须通过
sessionFactory.createSession()获取新会话,防止出现会话状态冲突。
内容的提问来源于stack exchange,提问作者Frank Leverenz
相关产品推荐
相关产品推荐

