NiFi需求:为FlowFile添加ContentDup属性标记是否含重复行
NiFi实现文件重复行检测并添加标记属性方案
直接用ExecuteScript处理器结合Groovy脚本即可实现需求,无需拆分FlowFile,仅添加属性标记:
实现步骤
- 添加ExecuteScript处理器,在"Script Language"选项中选择
Groovy - 替换默认脚本为以下代码:
import org.apache.nifi.processor.io.StreamCallback import java.nio.charset.StandardCharsets def flowFile = session.get() if (!flowFile) return flowFile = session.write(flowFile, { inputStream, outputStream -> // 读取文件所有行 def lines = inputStream.readLines(StandardCharsets.UTF_8) // 对比原行数与去重后行数判断是否有重复 def hasDuplicates = lines.size() != lines.unique().size() // 设置标记属性 flowFile = session.putAttribute(flowFile, "ContentDup", hasDuplicates ? "Yes" : "No") // 原内容原样写出,不做修改 outputStream.write(lines.join(System.lineSeparator()).getBytes(StandardCharsets.UTF_8)) } as StreamCallback) session.transfer(flowFile, REL_SUCCESS)
大文件适配优化
如果处理GB级大文件,避免全量读入内存溢出,改用流式逐行校验方案:
import org.apache.nifi.processor.io.StreamCallback import java.nio.charset.StandardCharsets import java.util.HashSet def flowFile = session.get() if (!flowFile) return def hasDuplicates = false def seenLines = new HashSet<String>() flowFile = session.write(flowFile, { inputStream, outputStream -> def reader = new BufferedReader(new InputStreamReader(inputStream, StandardCharsets.UTF_8)) def writer = new BufferedWriter(new OutputStreamWriter(outputStream, StandardCharsets.UTF_8)) String line while ((line = reader.readLine()) != null) { // 若该行已存在则标记重复 if (!seenLines.add(line)) { hasDuplicates = true } // 原样写出内容 writer.write(line) writer.newLine() } writer.flush() } as StreamCallback) flowFile = session.putAttribute(flowFile, "ContentDup", hasDuplicates ? "Yes" : "No") session.transfer(flowFile, REL_SUCCESS)
内容的提问来源于stack exchange,提问作者user21803274
相关产品推荐
相关产品推荐

