如何在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
相关产品推荐
相关产品推荐

