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

NiFi处理器onTrigger中添加的FlowFile属性无法即时查看的问题

FlowFile属性在Processor的onTrigger方法内不可见问题解析

环境

  • NiFi版本:1.23.2
  • JDK:OpenJDK 11.0.12
  • 操作系统:Ubuntu 20.04 64位LTS或Windows 11专业版

问题描述

在两处为FlowFile添加属性:

  1. 在TestSampleProcessor.java中模拟FlowFile已有属性
  2. 在SampleProcessor.java的onTrigger(...)方法中添加IN_SIDE_ATTR_1和IN_SIDE_ATTR_2

从输出日志可见,onTrigger(...)内的日志无法找到这两个新增属性,但TestSampleProcessor的日志中能正常显示。

相关代码

SampleProcessor.java

import org.apache.nifi.annotation.documentation.CapabilityDescription;
import org.apache.nifi.components.PropertyDescriptor;
import org.apache.nifi.flowfile.FlowFile;
import org.apache.nifi.logging.ComponentLog;
import org.apache.nifi.processor.*;

@CapabilityDescription("Sample processor for testing log message")
class SampleProcessor extends AbstractProcessor {
    public static final Relationship REL_SUCCESS = new Relationship.Builder()
            .name("success")
            .description("Test")
            .build();
    private List<PropertyDescriptor> descriptors = new ArrayList<>();

    private Set<Relationship> relationships = new HashSet<>();
    private ComponentLog log;

    @Override
    protected void init(final ProcessorInitializationContext context) {

        log = getLogger();
        //  relationships =
        relationships.add(REL_SUCCESS);
        relationships = Collections.unmodifiableSet(relationships);
    }

    @Override
    public Set<Relationship> getRelationships() {
        return this.relationships;
    }


    @Override
    public void onTrigger(final ProcessContext context, final ProcessSession session) {
        FlowFile flowFile = session.get();
        if (flowFile == null) return;
        InputStream inputStream = session.read(flowFile);
        BufferedReader bufferedReader = new BufferedReader(new InputStreamReader(inputStream, "utf-8"));
        String content = "", line = null;
        while ((line = bufferedReader.readLine())!= null) {
            content += line;
        }
        bufferedReader.close();
        println("log.isDebugEnabled():"+ log.isDebugEnabled());
        log.info("BEGIN to TEST XXXXXXXXXXXXXXXXXXXXX");
        log.debug("TEST =================================");
        log.debug("content is {}", content);
        log.debug("1. attribute count {}", flowFile.getAttributes().size());
        session.putAttribute(flowFile, "IN_SIDE_ATTR_1", "ABCD");
        session.putAttribute(flowFile, "IN_SIDE_ATTR_2", "EFGHI");
        log.debug("2. attribute count {}", flowFile.getAttributes().size());
        flowFile.getAttributes().keySet().stream().forEach (key ->{
            String value = flowFile.getAttribute(key);
            log.debug(String.format("In processor: [%s]=[%s]", key, value));
        });
        session.transfer(flowFile, REL_SUCCESS);
    }
}

TestSampleProcessor.java

import groovy.util.logging.Slf4j;
import org.apache.nifi.util.MockFlowFile;
import org.apache.nifi.util.TestRunner;
import org.apache.nifi.util.TestRunners;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;

@Slf4j
class TestSampleProcessor {
    private TestRunner testRunner;
    @BeforeEach
    public void init() {
        testRunner = TestRunners.newTestRunner(SampleProcessor.class);
    }
    @Test
    public void testAddAttribute() {
        String lines = "HELLO, WORLD";
        Map<String, String> attributes = new HashMap<>();
        attributes.put("ATTR_1_ADD_FROM_OUTSIDE", UUID.randomUUID().toString());
        attributes.put("ATTR_2_ADD_FROM_OUTSIDE", "/Retry");

        testRunner.enqueue(lines, attributes);
        testRunner.run();
        testRunner.assertAllFlowFilesTransferred(SampleProcessor.REL_SUCCESS, 1);
        MockFlowFile flowFile = testRunner.getFlowFilesForRelationship(SampleProcessor.REL_SUCCESS).get(0);
        log.info("Check all attributes");
        flowFile.attributes.keySet().forEach (key-> {
            log.debug(String.format("DEBUG. [%s]=[%s]", key, flowFile.getAttribute(key)));
        });
    }
}

输出日志

log.isDebugEnabled():true
2023-10-11 09:14:46:895 +0800 INFO SampleProcessor - SampleProcessor[id=39f7be77-f23e-429b-8b59-7a165e3de2ed] BEGIN to TEST XXXXXXXXXXXXXXXXXXXXX
2023-10-11 09:14:46:896 +0800 DEBUG SampleProcessor - SampleProcessor[id=39f7be77-f23e-429b-8b59-7a165e3de2ed] TEST =================================
2023-10-11 09:14:46:897 +0800 DEBUG SampleProcessor - SampleProcessor[id=39f7be77-f23e-429b-8b59-7a165e3de2ed] content is HELLO, WORLD
2023-10-11 09:14:46:906 +0800 DEBUG SampleProcessor - SampleProcessor[id=39f7be77-f23e-429b-8b59-7a165e3de2ed] 1. attribute count 5
2023-10-11 09:14:46:906 +0800 DEBUG SampleProcessor - SampleProcessor[id=39f7be77-f23e-429b-8b59-7a165e3de2ed] 2. attribute count 5
2023-10-11 09:14:46:928 +0800 DEBUG SampleProcessor - SampleProcessor[id=39f7be77-f23e-429b-8b59-7a165e3de2ed] In processor: [filename]=[88495457086800.mockFlowFile]
2023-10-11 09:14:46:928 +0800 DEBUG SampleProcessor - SampleProcessor[id=39f7be77-f23e-429b-8b59-7a165e3de2ed] In processor: [path]=[target]
2023-10-11 09:14:46:929 +0800 DEBUG SampleProcessor - SampleProcessor[id=39f7be77-f23e-429b-8b59-7a165e3de2ed] In processor: [uuid]=[7a5fbff1-9d23-43a6-93b6-fd73d8d4974d]
2023-10-11 09:14:46:929 +0800 DEBUG SampleProcessor - SampleProcessor[id=39f7be77-f23e-429b-8b59-7a165e3de2ed] In processor: [ATTR_1_ADD_FROM_OUTSIDE]=[4eac4e85-362f-48f5-806c-33e322104cb1]
2023-10-11 09:14:46:930 +0800 DEBUG SampleProcessor - SampleProcessor[id=39f7be77-f23e-429b-8b59-7a165e3de2ed] In processor: [ATTR_2_ADD_FROM_OUTSIDE]=[/Retry]
2023-10-11 09:14:46:941 +0800 INFO TestSampleProcessor - Check all attributes
2023-10-11 09:14:46:944 +0800 DEBUG TestSampleProcessor - DEBUG. [filename]=[88495457086800.mockFlowFile]
2023-10-11 09:14:46:944 +0800 DEBUG TestSampleProcessor - DEBUG. [path]=[target]
2023-10-11 09:14:46:944 +0800 DEBUG TestSampleProcessor - DEBUG. [uuid]=[7a5fbff1-9d23-43a6-93b6-fd73d8d4974d]
2023-10-11 09:14:46:944 +0800 DEBUG TestSampleProcessor - DEBUG. [ATTR_1_ADD_FROM_OUTSIDE]=[4eac4e85-362f-48f5-806c-33e322104cb1]
2023-10-11 09:14:46:944 +0800 DEBUG TestSampleProcessor - DEBUG. [ATTR_2_ADD_FROM_OUTSIDE]=[/Retry]
2023-10-11 09:14:46:944 +0800 DEBUG TestSampleProcessor - DEBUG. [IN_SIDE_ATTR_1]=[ABCD]
2023-10-11 09:14:46:944 +0800 DEBUG TestSampleProcessor - DEBUG. [IN_SIDE_ATTR_2]=[EFGHI]

原因分析

NiFi中的FlowFile是不可变对象,调用session.putAttribute()方法时,并不会修改原来的flowFile实例,而是会创建一个包含新属性的FlowFile新实例。你在代码中直接使用原来的flowFile对象去获取属性,自然看不到新增的属性。

而测试代码中能看到属性,是因为testRunner.run()执行完成后,会将所有会话中的修改提交,最终拿到的是包含所有属性的最终FlowFile实例。

解决方法

在调用session.putAttribute()后,需要接收返回的新FlowFile实例,并使用这个新实例进行后续操作。修改onTrigger方法的关键代码如下:

// 替换原来的putAttribute调用,接收返回的新FlowFile
flowFile = session.putAttribute(flowFile, "IN_SIDE_ATTR_1", "ABCD");
flowFile = session.putAttribute(flowFile, "IN_SIDE_ATTR_2", "EFGHI");

// 此时使用新的flowFile实例获取属性,就能看到新增的属性了
log.debug("2. attribute count {}", flowFile.getAttributes().size());
flowFile.getAttributes().keySet().stream().forEach (key ->{
    String value = flowFile.getAttribute(key);
    log.debug(String.format("In processor: [%s]=[%s]", key, value));
});

另外,代码中读取FlowFile内容时,没有正确处理流关闭的异常风险,建议使用try-with-resources语法自动关闭流,避免资源泄漏:

try (InputStream inputStream = session.read(flowFile);
     BufferedReader bufferedReader = new BufferedReader(new InputStreamReader(inputStream, "utf-8"))) {
    String content = "", line = null;
    while ((line = bufferedReader.readLine())!= null) {
        content += line;
    }
    // 后续处理逻辑
} catch (IOException e) {
    log.error("Failed to read FlowFile content", e);
    // 建议新增失败关系并处理异常场景
    session.transfer(flowFile, REL_FAILURE);
    return;
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 15:17:02