NiFi处理器onTrigger中添加的FlowFile属性无法即时查看的问题
FlowFile属性在Processor的onTrigger方法内不可见问题解析
环境
- NiFi版本:1.23.2
- JDK:OpenJDK 11.0.12
- 操作系统:Ubuntu 20.04 64位LTS或Windows 11专业版
问题描述
在两处为FlowFile添加属性:
- 在
TestSampleProcessor.java中模拟FlowFile已有属性 - 在
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
相关产品推荐
相关产品推荐

