如何在Apache NiFi中对比两个FlowFile内容并提取差异
在Apache NiFi中实现FlowFile逐行对比并提取差异
需求说明
需要逐行对比两个来自数据库连接的FlowFile内容,提取差异数据并生成带变更说明的结果。两个FlowFile的内容格式(以~为分隔符的键值对)及期望输出如下:
第一个FlowFile内容
123|3~abc|123|Hyderabad 456|3~abc|456|delhi 789|3~abc|789|delhi
第二个FlowFile内容
1234|3~abc|123|Hyderabad 456|3~abc|456|delhi 789|3~abc|789|chennai
期望差异输出
1234|3~abc|123|Hyderabad(123 is changed to 1234) 789|3~abc|789|chennai(delhi is changed to chennai)
数据量可达数千行,以下是可行的实现方案:
内置处理器组合方案
可以通过以下内置处理器组合实现,无需自定义脚本,但配置相对繁琐:
- MergeContent:将两个FlowFile合并为一个,可选择"Concatenate"模式,在两个文件内容间添加标识行区分,确保后续能准确拆分对应行。
- SplitText:将合并后的FlowFile按行拆分,每行生成独立的FlowFile。
- PartitionRecord:以每行的唯一标识(比如示例中
~分隔后的第二部分abc|xxx)作为分区键,让两个原FlowFile中对应键的行进入同一分区。 - MergeRecord:将同一分区内的两行(分别来自两个原FlowFile)合并为单条记录,便于对比。
- UpdateRecord + ReplaceText:通过UpdateRecord标记差异字段,再用ReplaceText拼接成带说明的差异行;但这种方式仅适用于简单的差异场景,复杂的多字段对比会增加配置复杂度。
Groovy脚本方案(高效灵活)
针对数千行的数据量,使用Groovy脚本配合ExecuteScript处理器是更直接高效的选择,核心实现步骤如下:
- 先通过MergeContent将两个FlowFile合并,或确保它们能被同一个ExecuteScript实例读取。
- 在Groovy脚本中完成行匹配、差异对比与结果生成:
- 读取两个FlowFile的内容并按行拆分,建立第一个文件的行-键映射(以唯一标识为键)。
- 遍历第二个文件的每一行,匹配第一个文件中对应键的行,逐字段对比差异。
- 拼接带差异说明的字符串,写入新的FlowFile。
示例Groovy脚本核心逻辑:
def flowFile1 = session.get() def flowFile2 = session.get() if (!flowFile1 || !flowFile2) { return } // 读取两个文件的内容并按行拆分 def content1 = flowFile1.read().text.split('\n') def content2 = flowFile2.read().text.split('\n') // 构建第一个文件的行映射,以~分隔的第二部分为唯一键 def lineMap = [:] content1.each { line -> def key = line.split('~')[1] lineMap[key] = line } def resultBuilder = new StringBuilder() content2.each { line2 -> def key = line2.split('~')[1] def line1 = lineMap[key] if (!line1) { resultBuilder.append(" ${line2}(new line)\n") } else { def parts1 = line1.split('\\|') def parts2 = line2.split('\\|') def diffDetails = [] // 逐字段对比差异 for (int i = 0; i < parts1.size(); i++) { if (parts1[i] != parts2[i]) { diffDetails.add("${parts1[i]} is changed to ${parts2[i]}") } } if (diffDetails.size() > 0) { resultBuilder.append(" ${line2}(${diffDetails.join(', ')})\n") } } } // 生成结果FlowFile def outFlowFile = session.create() outFlowFile = session.write(outFlowFile, { outputStream -> outputStream.write(resultBuilder.toString().getBytes('UTF-8')) } as OutputStreamCallback) session.transfer(outFlowFile, REL_SUCCESS) session.remove(flowFile1) session.remove(flowFile2)
方案选择建议
- 若偏好使用内置组件且数据量较小,可尝试处理器组合方案,但需注意配置细节以保证行匹配准确性。
- 对于数千行的规模,Groovy脚本方案更简洁高效,能灵活处理多字段差异、新增行等复杂场景。
内容的提问来源于stack exchange,提问作者Amarnatha Reddy
相关产品推荐
相关产品推荐

