如何拆分flowfile属性中的ids值生成多个新flowfile解决InvokeHTTP报错
NiFi FlowFile ids属性二分拆分重试方案
场景说明
现有无内容FlowFile携带ids属性,属性值为逗号分隔的id列表,示例如下:
1018866556,1018878837,1018522766,1018522773,1018522788,1018522790,1018522797, 1018522959,1018522963,1018522968,1018522972,1018522981,1018511143,1018511174
在InvokeHTTP中使用该属性时,因数据量过大触发报错,需要实现失败自动二分拆分,拆分后重试的逻辑。
整体流程设计
- 初始FlowFile进入链路后,先通过UpdateAttribute新增两个控制属性:
max_split_times:最大拆分次数,初始值建议设为10,避免非参数长度导致的错误触发无限拆分current_split_level:当前拆分层级,初始值设为0,每次拆分后+1
- 配置InvokeHTTP处理器,
success关系直接接入后续业务流程,failure、RETRY、NO_RETRY(根据实际错误类型调整)关系接入拆分处理节点 - 拆分完成后的FlowFile重新回流到InvokeHTTP前的节点重试
核心拆分逻辑实现
使用ExecuteScript处理器完成拆分,脚本引擎选择Groovy,参考脚本如下:
def flowFile = session.get() if (!flowFile) return try { // 读取ids属性,分割为id列表,过滤空值 def idsStr = flowFile.getAttribute('ids')?.trim() if (!idsStr) { session.transfer(flowFile, REL_FAILURE) return } def ids = idsStr.split(',').collect { it.trim() }.findAll { it } def idCount = ids.size() // 达到拆分上限或无法继续拆分,直接路由到失败 def maxSplit = flowFile.getAttribute('max_split_times')?.toInteger() ?: 10 def currentLevel = flowFile.getAttribute('current_split_level')?.toInteger() ?: 0 if (idCount <= 1 || currentLevel >= maxSplit) { session.transfer(flowFile, REL_FAILURE) return } // 二分拆分 def mid = (idCount / 2) as int def firstHalf = ids.subList(0, mid).join(',') def secondHalf = ids.subList(mid, idCount).join(',') // 生成第一个拆分后的FlowFile def firstFlowFile = session.create(flowFile) firstFlowFile = session.putAllAttributes(firstFlowFile, flowFile.getAttributes()) firstFlowFile = session.putAttribute(firstFlowFile, 'ids', firstHalf) firstFlowFile = session.putAttribute(firstFlowFile, 'current_split_level', (currentLevel + 1).toString()) session.transfer(firstFlowFile, REL_SUCCESS) // 生成第二个拆分后的FlowFile def secondFlowFile = session.create(flowFile) secondFlowFile = session.putAllAttributes(secondFlowFile, flowFile.getAttributes()) secondFlowFile = session.putAttribute(secondFlowFile, 'ids', secondHalf) secondFlowFile = session.putAttribute(secondFlowFile, 'current_split_level', (currentLevel + 1).toString()) session.transfer(secondFlowFile, REL_SUCCESS) // 删除原始FlowFile session.remove(flowFile) } catch (Exception e) { session.transfer(flowFile, REL_FAILURE) }
额外配置说明
- 若InvokeHTTP将ids放在请求URL参数中,拆分后需确认属性值长度未超过服务端限制的最大URL长度,一般建议拆分后单批id数量控制在200以内
- 若ids需要放在请求体中,可以在InvokeHTTP前新增ReplaceText处理器,将
ids属性值写入FlowFile内容,拆分时不需要修改内容,仅修改属性即可 - 拆分后重试的频率可以通过在拆分链路中加入ControlRate处理器控制,避免短时间内大量请求冲击接口
内容的提问来源于stack exchange,提问作者miroha
相关产品推荐
相关产品推荐

