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

Apache NiFi自定义处理器并发任务Java实现技术咨询

NiFi Concurrent Tasks 机制说明

核心原理

NiFi调度框架会根据你配置的Concurrent Tasks数量,创建对应数目的工作线程,每个线程独立调用Processor的onTrigger(ProcessContext context, ProcessSession session)方法。这意味着自定义Processor必须保证线程安全——不能使用非线程安全的成员变量,所有状态要么依托ProcessSession管理,要么采用JUC包下的线程安全容器(如ConcurrentHashMap、AtomicInteger)。

自定义Processor适配要点

  • 禁止使用普通HashMap、ArrayList这类非线程安全的实例变量,若需维护状态,优先选用线程安全类。
  • 不要在onTrigger方法中加全局锁(除非业务强制要求),否则会完全抵消并发调度的效果。
  • 依赖的外部资源(如数据库连接池)需本身支持并发访问。

代码示例

以下是一个线程安全的自定义Processor示例,演示如何适配并发任务调度:

import org.apache.nifi.annotation.behavior.InputRequirement;
import org.apache.nifi.annotation.documentation.CapabilityDescription;
import org.apache.nifi.annotation.documentation.Tags;
import org.apache.nifi.flowfile.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 java.util.Collections;
import java.util.HashSet;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicInteger;

@Tags({"concurrent", "demo"})
@CapabilityDescription("支持并发任务的自定义Processor示例,展示线程安全的状态管理方式")
@InputRequirement(InputRequirement.Requirement.INPUT_REQUIRED)
public class ConcurrentSafeProcessor extends AbstractProcessor {

    // 线程安全计数器,统计已处理FlowFile总数
    private final AtomicInteger processedTotal = new AtomicInteger(0);
    // 线程安全集合,存储已处理FlowFile的ID与文件名映射
    private final ConcurrentHashMap<String, String> processedFlowFileMap = new ConcurrentHashMap<>();

    public static final Relationship REL_SUCCESS = new Relationship.Builder()
            .name("success")
            .description("处理完成的FlowFile流向此关系")
            .build();

    private static final Set<Relationship> RELATIONSHIPS;
    static {
        final Set<Relationship> relationships = new HashSet<>();
        relationships.add(REL_SUCCESS);
        RELATIONSHIPS = Collections.unmodifiableSet(relationships);
    }

    @Override
    public Set<Relationship> getRelationships() {
        return RELATIONSHIPS;
    }

    @Override
    public void onTrigger(final ProcessContext context, final ProcessSession session) throws ProcessException {
        FlowFile flowFile = session.get();
        if (flowFile == null) {
            return;
        }

        try {
            // 模拟业务处理耗时操作
            Thread.sleep(150);

            // 线程安全更新状态
            int currentTotal = processedTotal.incrementAndGet();
            processedFlowFileMap.put(flowFile.getId(), flowFile.getAttribute("filename"));

            getLogger().info("累计处理FlowFile:{},当前处理文件:{}", currentTotal, flowFile.getAttribute("filename"));

            // 转移FlowFile至成功关系
            session.transfer(flowFile, REL_SUCCESS);
        } catch (InterruptedException e) {
            getLogger().error("处理FlowFile {}失败", flowFile.getId(), e);
            session.rollback();
            Thread.currentThread().interrupt();
        }
    }
}

关键细节说明

  • 线程安全状态:用AtomicInteger实现无锁的计数操作,ConcurrentHashMap存储业务状态,避免多线程下的竞态条件。
  • 无全局锁设计:每个线程独立处理分配到的FlowFile,最大化并发调度的效率。
  • 异常处理:捕获中断异常后恢复线程中断状态,确保NiFi调度框架能正确管理线程生命周期。

内容的提问来源于stack exchange,提问作者Raghav07

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 09:12:43