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

如何在Apache NiFi中用Groovy脚本逐行对比两个CSV文件

Apache NiFi中对比Flow File与数据库CSV数据的Groovy脚本实现

场景说明

在Apache NiFi流程中,Flow File携带一份CSV数据,数据库中存储着同结构但存在差异的CSV数据,需要通过Groovy脚本完成以下操作:

  • 识别Flow File中数据库不存在的记录(需新增)
  • 识别Flow File中与数据库同主键但字段有差异的记录(需更新)

完整Groovy脚本

@Grab('org.apache.commons:commons-csv:1.9.0')
import org.apache.commons.csv.CSVFormat
import org.apache.commons.csv.CSVParser
import org.apache.commons.csv.CSVRecord
import java.sql.Connection
import java.sql.DriverManager
import java.sql.ResultSet
import java.sql.Statement

// 1. 读取Flow File中的CSV内容
def flowFileContent = new String(flowFile.getContent())
def csvFormat = CSVFormat.DEFAULT.withHeader()
def flowFileRecords = csvFormat.parse(new StringReader(flowFileContent)).getRecords()

// 2. 从数据库读取目标表数据
def dbConfig = [
    url: 'jdbc:mysql://your-db-host:3306/your-db',
    user: 'db-user',
    password: 'db-pass',
    table: 'your-target-table',
    primaryKey: 'id' // 替换为实际主键字段名
]

Connection conn = null
Statement stmt = null
ResultSet rs = null
def dbRecordsMap = [:]

try {
    Class.forName('com.mysql.cj.jdbc.Driver')
    conn = DriverManager.getConnection(dbConfig.url, dbConfig.user, dbConfig.password)
    stmt = conn.createStatement()
    def query = "SELECT * FROM ${dbConfig.table}"
    rs = stmt.executeQuery(query)
    
    // 把数据库记录转成主键为key的Map,方便快速匹配
    def columnCount = rs.getMetaData().getColumnCount()
    while (rs.next()) {
        def recordMap = [:]
        for (int i = 1; i <= columnCount; i++) {
            def colName = rs.getMetaData().getColumnName(i)
            recordMap[colName] = rs.getString(colName)
        }
        dbRecordsMap[recordMap[dbConfig.primaryKey]] = recordMap
    }
} finally {
    // 确保资源释放,避免泄漏
    rs?.close()
    stmt?.close()
    conn?.close()
}

// 3. 对比数据,筛选新增、更新记录
def insertRecords = []
def updateRecords = []

flowFileRecords.each { flowRecord ->
    def pkValue = flowRecord.get(dbConfig.primaryKey)
    def dbRecord = dbRecordsMap.get(pkValue)
    
    if (!dbRecord) {
        // 数据库无此主键,标记为新增
        insertRecords.add(flowRecord.toMap())
    } else {
        // 逐字段对比,判断是否有变更
        def hasChange = false
        flowRecord.toMap().each { key, value ->
            if (dbRecord[key] != value) {
                hasChange = true
                return // 跳出循环
            }
        }
        if (hasChange) {
            updateRecords.add(flowRecord.toMap())
        }
    }
}

// 4. 输出结果(可根据需求调整格式)
def result = [
    "insert_count": insertRecords.size(),
    "update_count": updateRecords.size(),
    "insert_records": insertRecords,
    "update_records": updateRecords
]

// 将结果转为JSON写入Flow File
flowFile = session.write(flowFile, { outputStream ->
    outputStream.write(new groovy.json.JsonBuilder(result).toPrettyString().getBytes('UTF-8'))
} as OutputStreamCallback)

// 可选:把统计数写入Flow File属性
session.putAttribute(flowFile, 'insert.count', insertRecords.size().toString())
session.putAttribute(flowFile, 'update.count', updateRecords.size().toString())

session.transfer(flowFile, REL_SUCCESS)

关键注意点

  • 主键指定:必须明确唯一主键字段,这是匹配记录的核心依据
  • 依赖引入:通过@Grab自动引入Apache Commons CSV库,用于解析CSV内容
  • 资源管理:数据库连接、结果集等资源必须在finally块中关闭,避免资源泄漏
  • 字段对比:脚本会逐字段对比所有值,只要有一个字段差异就标记为更新,可根据需求修改对比逻辑

内容的提问来源于stack exchange,提问作者Amarnatha Reddy

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 00:33:13