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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 03:48:21