Apache NiFi中如何按相同列值合并多行数据为单行?
解决方案
步骤1:拆分每行为独立FlowFile
使用SplitText处理器,配置如下:
- Line Split Count: 1
- Header Line Count: 0(你的数据无表头)
原FlowFile会被拆分为3个单独的FlowFile,每个对应一行原始数据。
步骤2:提取分组键与目标字段
使用ExtractText处理器,为每个拆分后的FlowFile提取两个属性:
- 分组键(第二列值):
- 属性名:
group_key - 正则表达式:
^\S+\t(\S+)\t\S+$(匹配制表符分隔的第二列)
- 属性名:
- 需要合并的字段(第三列值):
- 属性名:
content_part - 正则表达式:
^\S+\t\S+\t(\S+)$(匹配制表符分隔的第三列)
- 属性名:
步骤3:按分组键合并内容
使用MergeContent处理器,按group_key属性分组合并:
- Merge Strategy:
Bin-Packing Algorithm - Attribute Name for Correlation:
group_key - Delimiter Strategy:
Text - Demarcator:
(空格,用于分隔合并的第三列内容) - Minimum Number of Entries: 1
- Maximum Number of Entries: 1000(可根据数据量调整)
这一步会将同一分组下的所有第三列值用空格拼接,得到AAA BBB CCC。
步骤4:生成最终输出格式
方法一:ReplaceText处理器
- Search Value:
(.*) - Replacement Value:
${group_key},$1 - Replacement Strategy:
Replace Entire Text
直接将合并后的内容替换为101,AAA BBB CCC格式。
方法二:ExecuteScript处理器(Groovy脚本)
若需更灵活的处理,可使用Groovy脚本一次性组装结果:
def flowFile = session.get() if (!flowFile) return def groupKey = flowFile.getAttribute('group_key') def mergedContent = flowFile.read().text.trim() def finalContent = "${groupKey},${mergedContent}" flowFile = session.write(flowFile, { out -> out.write(finalContent.getBytes('UTF-8')) } as OutputStreamCallback) session.transfer(flowFile, REL_SUCCESS)
替代方案:ExecuteScript直接处理整文件
无需拆分合并,用Groovy脚本直接读取整个FlowFile内容并处理:
def flowFile = session.get() if (!flowFile) return def content = flowFile.read().text.split('\n').collect { line -> def parts = line.split('\t') [parts[1], parts[2]] // 提取第二列和第三列 } def groupKey = content[0][0] def mergedParts = content.collect { it[1] }.join(' ') def finalContent = "${groupKey},${mergedParts}" flowFile = session.write(flowFile, { out -> out.write(finalContent.getBytes('UTF-8')) } as OutputStreamCallback) session.transfer(flowFile, REL_SUCCESS)
内容的提问来源于stack exchange,提问作者Amarnatha Reddy
相关产品推荐
相关产品推荐

