NiFi中汇总多个FlowFile的source_count值的方法咨询
NiFi中汇总多个FlowFile的source_count值的方法咨询
嗨,针对你在NiFi里要汇总多个FlowFile中source_count值的需求,我给你整理了几个实用的方法,你可以根据自己的场景灵活选择:
方法一:MergeContent + ExecuteGroovyScript(适合批量小文件场景)
- 第一步:用
MergeContent处理器把所有包含source_count的FlowFile合并到一起。配置时可以选择Concatenate模式,让每个原FlowFile的JSON内容单独占一行,这样后续解析更方便。记得设置好合并的触发条件(比如达到指定数量或大小)。 - 第二步:用
ExecuteGroovyScript处理器编写脚本,读取合并后的内容,逐个解析JSON里的source_count并累加,最后输出包含总数的新FlowFile。示例脚本如下:
def flowFile = session.get() if (!flowFile) return def totalCount = 0 // 逐行读取合并后的内容 flowFile.read().withReader { reader -> reader.eachLine { line -> def jsonData = new groovy.json.JsonSlurper().parseText(line) totalCount += jsonData.source_count as Integer } } // 把总数写入新的JSON内容 def resultContent = new ByteArrayInputStream("{\"total_source_count\":${totalCount}}".getBytes('UTF-8')) flowFile = session.write(flowFile, { outputStream -> outputStream << resultContent } as OutputStreamCallback) session.transfer(flowFile, REL_SUCCESS)
方法二:DistributedMapCache 分布式累加(适合集群场景)
如果你的NiFi是集群部署,用这个方法能保证跨节点的累加一致性:
- 第一步:用
EvaluateJsonPath处理器把每个FlowFile里的source_count提取为属性,比如设置Destination为Attribute,Json Path表达式为$.source_count,属性名设为current_count。 - 第二步:用
DistributedMapCachePut处理器,配置缓存的Key为固定值(比如total_source_counter),Value设置为表达式${cache.value:plus(${current_count})},这样每来一个FlowFile就自动把当前值加到缓存的总数里。 - 第三步:当所有FlowFile处理完成后,用
DistributedMapCacheGet处理器获取缓存里的总数,生成包含结果的FlowFile即可。
方法三:JoltTransformJSON + MergeContent + 表达式语言(无脚本方案)
如果你不想写代码,可以用纯处理器组合实现:
- 第一步:用
JoltTransformJSON把每个FlowFile的JSON转换为纯数字格式,Jolt规格可以写:
这样就能把[ { "operation": "modify-overwrite-beta", "spec": { "source_count": "=toInteger" } }, { "operation": "remove", "spec": { "*": "" } } ]{"source_count":"200"}转成200。 - 第二步:用
MergeContent把所有数字合并为逗号分隔的字符串(比如200,150,300)。 - 第三步:用
UpdateAttribute或者ReplaceText处理器,用NiFi表达式语言计算总和:${merged_content:split(','):sum()},最后把结果生成新的JSON或者属性。
注意事项
- 如果处理的FlowFile数量极大,MergeContent要合理设置批次大小,避免内存占用过高。
- 用表达式语言处理时,要确保所有
source_count都是有效的数字格式,避免计算出错。
备注:内容来源于stack exchange,提问作者Ram
相关产品推荐
相关产品推荐

