如何在Apache NiFi中调用Java程序或Jar包?附具体实现需求
实现Apache NiFi集成TestDemo逻辑的方案
方案一:自定义NiFi处理器(推荐,贴合NiFi数据流模型)
1. 重构TestDemo代码结构
把核心业务逻辑从main方法中抽离为可复用的工具类,解耦入口逻辑与业务逻辑:
public class CSVCompareUtils { // 保留原读取逻辑,新增InputStream适配方法(适配NiFi数据流) public static Map<String, String> readCSV(String filePath) throws IOException { // 原readCSV实现代码 } public static Map<String, String> readCSVFromStream(InputStream in) throws IOException { // 基于InputStream实现文件内容读取,适配NiFi FlowFile } public static void compareAndMergeCSVs(Map<String, String> file1Data, Map<String, String> file2Data) { // 原比较合并逻辑 } }
2. 编写自定义NiFi处理器
继承AbstractProcessor,实现核心处理方法onTrigger,对接NiFi的FlowFile流转机制:
import org.apache.nifi.processor.AbstractProcessor; import org.apache.nifi.processor.ProcessContext; import org.apache.nifi.processor.ProcessSession; import org.apache.nifi.processor.Relationship; import org.apache.nifi.processor.exception.ProcessException; import org.apache.nifi.processor.io.InputStreamCallback; import java.io.InputStream; import java.util.Map; public class CSVCompareProcessor extends AbstractProcessor { // 定义处理器输出关系:成功、失败 public static final Relationship REL_SUCCESS = new Relationship.Builder() .name("success") .description("处理完成的FlowFile") .build(); public static final Relationship REL_FAILURE = new Relationship.Builder() .name("failure") .description("处理失败的FlowFile") .build(); @Override public void onTrigger(ProcessContext context, ProcessSession session) throws ProcessException { FlowFile flowFile = session.get(); if (flowFile == null) { return; } try { // 从处理器配置中获取第二个文件路径 String secondFilePath = context.getProperty("SECOND_FILE_PATH").getValue(); // 读取当前FlowFile的内容(第一个文件数据) Map<String, String> file1Data = CSVCompareUtils.readCSVFromStream(session.read(flowFile)); // 读取第二个文件数据 Map<String, String> file2Data = CSVCompareUtils.readCSV(secondFilePath); // 执行比较合并逻辑 CSVCompareUtils.compareAndMergeCSVs(file1Data, file2Data); // 流转成功的FlowFile session.transfer(flowFile, REL_SUCCESS); } catch (Exception e) { getLogger().error("CSV文件处理失败", e); // 流转失败的FlowFile session.transfer(flowFile, REL_FAILURE); } } }
3. 打包部署
- 将自定义处理器与
CSVCompareUtils打包成Jar包,排除NiFi核心依赖(如nifi-api、nifi-framework-api,避免版本冲突)。 - 将Jar包放入NiFi节点的
lib目录,重启NiFi服务。 - 在NiFi UI中找到自定义处理器,拖入画布后配置
SECOND_FILE_PATH等属性,连接上下游组件即可使用。
方案二:使用NiFi原生处理器调用Jar包
1. 打包TestDemo为可执行Jar
确保Jar包包含所有依赖,Manifest文件中指定主类为TestDemo。
2. 使用ExecuteProcess处理器
- 配置
Command为java,Arguments为-jar /绝对路径/testdemo.jar。 - 说明:此方式调用外部进程,需NiFi服务器具备Java环境,适合处理本地固定路径的文件,但无法直接对接NiFi数据流。
3. 使用ExecuteScript处理器(更灵活)
通过Groovy脚本调用Jar包中的工具类,直接处理FlowFile:
import com.yourpackage.CSVCompareUtils import org.apache.nifi.processor.io.InputStreamCallback def flowFile = session.get() if (!flowFile) return try { def secondFilePath = context.getProperty("SECOND_FILE_PATH").value def file1Data = [:] // 读取FlowFile内容 session.read(flowFile, new InputStreamCallback() { void process(InputStream in) throws IOException { file1Data.putAll(CSVCompareUtils.readCSVFromStream(in)) } }) def file2Data = CSVCompareUtils.readCSV(secondFilePath) CSVCompareUtils.compareAndMergeCSVs(file1Data, file2Data) session.transfer(flowFile, REL_SUCCESS) } catch (e) { log.error("处理失败", e) session.transfer(flowFile, REL_FAILURE) }
- 将TestDemo的Jar包放入NiFi的
lib目录,或在ExecuteScript的Module Directory属性中指定Jar路径。
关键注意事项
- 数据流适配:优先使用FlowFile的
InputStream读取内容,避免直接读取本地路径,适配NiFi分布式流转特性。 - 依赖冲突:打包时排除NiFi自带的依赖类库,防止版本冲突导致处理器加载失败。
- 错误处理:必须捕获异常并流转失败的FlowFile,方便后续重试或异常链路处理。
内容的提问来源于stack exchange,提问作者Amarnatha Reddy
相关产品推荐
相关产品推荐

