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

如何使用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组合:

  1. MergeRecord:将所有FlowFile的Record合并为一个集合(Merge Strategy = Merge all records into a single RecordSet)
  2. ExecuteScript(Groovy脚本):提取所有字段名,固定前3列顺序,后续数字列按数值排序,将最终字段列表写入FlowFile属性all.columns
    import 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 00:20:52