如何使用Apache NiFi合并前3列相同、后续列不同的多份CSV文件?
用Apache NiFi合并结构相似但列数不同的CSV文件
需求场景
文件夹内的CSV文件满足以下特征:
- 前3列表头完全一致(示例中为
AAA,BBB,CCC) - 后续列的表头为数字,列数在2-11之间不等
- 需要合并所有文件为一个CSV,自动补全缺失列(值为
null),且仅保留一份完整表头
示例输入
文件1:
AAA,BBB,CCC,0,10,15 1,India,c,0,28,54 2,Taiwan,c,0,23,52 3,France,c,0,26,34 4,Japan,c,0,27,46
文件2:
AAA,BBB,CCC,0,5,15,30,40 1,Brazil,c,0,20,64,71,88 2,Russia,c,0,20,62,72,81 3,Poland,c,0,21,64,78,78 4,Litva,c,0,22,66,75,78
期望输出
AAA,BBB,CCC,0,5,10,15,30,40 1,India,c,0,null,28,54,null,null 2,Taiwan,c,0,null,23,52,null,null 3,France,c,0,null,26,34,null,null 4,Japan,c,0,null,27,46,null,null 1,Brazil,c,0,20,null,64,71,88 2,Russia,c,0,20,null,62,72,81 3,Poland,c,0,21,null,64,78,78 4,Litva,c,0,22,null,66,75,78
实现方案
直接使用Merge Content只能简单追加内容,无法处理列对齐和去重表头,需通过以下处理器组合实现:
1. 读取文件夹内的CSV文件
- 使用
ListFile:配置目标文件夹路径,按需设置递归扫描规则,定位所有CSV文件 - 使用
FetchFile:根据ListFile输出的文件路径,读取CSV内容生成FlowFile
2. 将CSV转换为结构化记录
使用ConvertCSVToRecord处理器:
- 配置
CSVReader:设置Header Line Count = 1,Schema Access Strategy = Infer Schema(自动识别表头和字段类型) - 输出的FlowFile会被转换成结构化Record格式,方便后续字段操作
3. 收集所有唯一表头字段
使用MergeRecord + ExecuteScript组合:
MergeRecord:将所有FlowFile的Record合并为一个集合(Merge Strategy = Merge all records into a single RecordSet)ExecuteScript(Groovy脚本):提取所有字段名,固定前3列顺序,后续数字列按数值排序,将最终字段列表写入FlowFile属性all.columnsimport org.apache.nifi.processor.io.StreamCallback import groovy.json.JsonSlurper import groovy.json.JsonBuilder def flowFile = session.get() if (!flowFile) return flowFile = session.write(flowFile, { inputStream, outputStream -> def records = new JsonSlurper().parse(inputStream) def fixedColumns = ['AAA', 'BBB', 'CCC'] def numericColumns = [] records.each { record -> record.keySet().each { key -> if (!fixedColumns.contains(key) && !numericColumns.contains(key)) { numericColumns.add(key) } } } // 按数值排序数字列,避免字符串排序导致的顺序错误 numericColumns.sort { it as Integer } def allColumns = fixedColumns + numericColumns flowFile = session.putAttribute(flowFile, 'all.columns', allColumns.join(',')) outputStream.write(new JsonBuilder(records).toByteArray()) } as StreamCallback) session.transfer(flowFile, REL_SUCCESS)
4. 统一所有记录的字段结构
使用ExecuteScript处理器,遍历每条Record,根据all.columns补充缺失字段(值设为null):
import org.apache.nifi.processor.io.StreamCallback import groovy.json.JsonSlurper import groovy.json.JsonBuilder def flowFile = session.get() if (!flowFile) return def allColumns = flowFile.getAttribute('all.columns').split(',') flowFile = session.write(flowFile, { inputStream, outputStream -> def records = new JsonSlurper().parse(inputStream) def normalizedRecords = records.collect { record -> def normalized = [:] allColumns.each { col -> normalized[col] = record[col] ?: 'null' } normalized } outputStream.write(new JsonBuilder(normalizedRecords).toByteArray()) } as StreamCallback) session.transfer(flowFile, REL_SUCCESS)
5. 生成最终合并后的CSV
使用MergeRecord处理器:
- 配置
CSVRecordSetWriter:设置Header Line Count = 1,Schema Access Strategy = Use Schema Text,将all.columns转换为Avro Schema(例如:{"type":"record","name":"merged","fields":[{"name":"AAA","type":"string"},{"name":"BBB","type":"string"},{"name":"CCC","type":"string"},{"name":"0","type":["string","null"]},...]}) - 设置
Merge Strategy = Merge all records into a single RecordSet,输出的FlowFile即为符合要求的合并CSV
注意事项
- 若文件数量较多,需调整
MergeRecord的Max Number of Records参数,避免内存溢出 ConvertCSVToRecord的Infer Schema需确保字段类型识别正确,若有特殊格式可自定义Schema- 数字列排序需按数值处理,避免字符串排序导致的顺序错误(如"10"排在"5"前面)
内容的提问来源于stack exchange,提问作者Bennet Turner
相关产品推荐
相关产品推荐

