NiFi开发求助:能否在@OnStopped注解中发送FlowFile?
Apache NiFi自定义处理器:@OnStopped中无法发送FlowFile的解决方案
嘿,我来帮你捋清楚这个问题:你不能在@OnStopped注解的方法里发送FlowFile或者操作ProcessSession。原因很直接:NiFi的处理器生命周期中,@OnStopped是处理器彻底停止前的最后阶段,此时ProcessSession已经被框架回收或者处于不可用状态,任何针对FlowFile的创建、属性修改或提交操作都会抛出异常,甚至导致处理器异常退出。
那怎么实现“处理器停止时生成带指定属性的FlowFile”这个需求呢?这里给你一个符合NiFi规范的可行方案:利用@OnStopping注解和一个状态标志位,在处理器收到停止信号后,在最后一次onTrigger执行时完成FlowFile的处理。
具体实现步骤
- 定义一个
volatile修饰的布尔标志,保证多线程下的可见性,用来标记处理器是否即将停止 - 在
@OnStopping方法中设置这个标志位(@OnStopping是处理器开始停止流程时调用,此时onTrigger可能还会执行最后一次) - 在
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
相关产品推荐
相关产品推荐

