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

NiFi开发求助:能否在@OnStopped注解中发送FlowFile?

Apache NiFi自定义处理器:@OnStopped中无法发送FlowFile的解决方案

嘿,我来帮你捋清楚这个问题:你不能在@OnStopped注解的方法里发送FlowFile或者操作ProcessSession。原因很直接:NiFi的处理器生命周期中,@OnStopped是处理器彻底停止前的最后阶段,此时ProcessSession已经被框架回收或者处于不可用状态,任何针对FlowFile的创建、属性修改或提交操作都会抛出异常,甚至导致处理器异常退出。

那怎么实现“处理器停止时生成带指定属性的FlowFile”这个需求呢?这里给你一个符合NiFi规范的可行方案:利用@OnStopping注解和一个状态标志位,在处理器收到停止信号后,在最后一次onTrigger执行时完成FlowFile的处理。

具体实现步骤

  1. 定义一个volatile修饰的布尔标志,保证多线程下的可见性,用来标记处理器是否即将停止
  2. 在@OnStopping方法中设置这个标志位(@OnStopping是处理器开始停止流程时调用,此时onTrigger可能还会执行最后一次)
  3. 在onTrigger方法中检查标志位,当标志位为true时,创建FlowFile并添加指定属性,最后提交会话

代码示例

import org.apache.nifi.annotation.behavior.OnStopping;
import org.apache.nifi.annotation.documentation.CapabilityDescription;
import org.apache.nifi.annotation.documentation.Tags;
import org.apache.nifi.components.PropertyDescriptor;
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.Set;

@Tags({"custom", "stop", "flowfile"})
@CapabilityDescription("处理器停止时生成带指定属性的FlowFile")
public class StopSignalProcessor extends AbstractProcessor {

    // 定义成功输出关系
    public static final Relationship REL_SUCCESS = new Relationship.Builder()
            .name("success")
            .description("携带停止信号的FlowFile发送到此关系")
            .build();

    // 停止标志位,volatile保证多线程环境下的可见性
    private volatile boolean isStopping = false;

    @Override
    public Set<Relationship> getRelationships() {
        return Collections.singleton(REL_SUCCESS);
    }

    @Override
    public Set<PropertyDescriptor> getSupportedPropertyDescriptors() {
        return Collections.emptySet();
    }

    @OnStopping
    public void onStopping(ProcessContext context) {
        // 收到停止信号,设置标志位
        isStopping = true;
        getLogger().info("处理器即将停止,将生成携带停止信号的FlowFile");
    }

    @Override
    public void onTrigger(final ProcessContext context, final ProcessSession session) throws ProcessException {
        // 检查是否处于停止流程中
        if (isStopping) {
            // 创建空FlowFile
            FlowFile flowFile = session.create();
            // 添加指定属性
            flowFile = session.putAttribute(flowFile, "ATTRIBUTE_SIGNAL", "PROCESSOR_STOPPED");
            // 发送到成功关系
            session.transfer(flowFile, REL_SUCCESS);
            // 提交会话,确保数据不丢失
            session.commit();
            // 让处理器不再被调度,避免重复执行
            context.yield();
            return;
        }

        // 这里是处理器正常运行时的业务逻辑,根据你的需求自行实现
        // ...
    }
}

关键注意事项

  • 绝对不要在@OnStopped中操作ProcessSession:这个阶段框架已经在清理资源,Session无法正常工作,强行操作会引发各种异常。
  • 用@OnStopping设置停止标志:@OnStopping是处理器停止流程的第一个回调,此时处理器还处于可调度状态,onTrigger会被执行最后一次,这是安全处理FlowFile的时机。
  • 务必提交Session:生成FlowFile后一定要调用session.commit(),否则数据会丢失。
  • 调用context.yield():处理完停止信号后,调用yield让处理器不再被调度,避免重复生成FlowFile。

这个方案完全符合NiFi的生命周期规范,能完美实现你想要的“处理器停止时生成带指定属性的FlowFile”的需求。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 09:46:05