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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.22 15:04:50