This is an automated email from the ASF dual-hosted git repository. joewitt pushed a commit to branch NIFI-16311 in repository https://gitbox.apache.org/repos/asf/nifi.git
commit 6f292c41e059789571c3a7b613b4e559c2db50af Author: Joseph Witt <[email protected]> AuthorDate: Mon Sep 7 13:57:10 2026 -0700 NIFI-16311 Initialize ParseSyslog parser in @OnScheduled CHARSET cannot change between triggers, so lazy construction in onTrigger was an unsynchronized check-then-act on the shared processor instance. Match ParseSyslog5424: build once at schedule time and publish the parser as volatile. Co-authored-by: Cursor <[email protected]> --- .../nifi/processors/standard/ParseSyslog.java | 15 +++++++-------- .../nifi/processors/standard/TestParseSyslog.java | 21 +++++++++++++++++++++ 2 files changed, 28 insertions(+), 8 deletions(-) diff --git a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/ParseSyslog.java b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/ParseSyslog.java index 0b4e30c86d4..e6e089e8e38 100644 --- a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/ParseSyslog.java +++ b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/ParseSyslog.java @@ -26,6 +26,7 @@ import org.apache.nifi.annotation.behavior.WritesAttributes; import org.apache.nifi.annotation.documentation.CapabilityDescription; import org.apache.nifi.annotation.documentation.SeeAlso; import org.apache.nifi.annotation.documentation.Tags; +import org.apache.nifi.annotation.lifecycle.OnScheduled; import org.apache.nifi.components.PropertyDescriptor; import org.apache.nifi.flowfile.FlowFile; import org.apache.nifi.processor.AbstractProcessor; @@ -91,7 +92,7 @@ public class ParseSyslog extends AbstractProcessor { REL_SUCCESS ); - private SyslogParser parser; + private volatile SyslogParser parser; @Override protected List<PropertyDescriptor> getSupportedPropertyDescriptors() { @@ -103,6 +104,11 @@ public class ParseSyslog extends AbstractProcessor { return RELATIONSHIPS; } + @OnScheduled + public void onScheduled(final ProcessContext context) { + parser = new SyslogParser(Charset.forName(context.getProperty(CHARSET).getValue())); + } + @Override public void onTrigger(final ProcessContext context, final ProcessSession session) throws ProcessException { FlowFile flowFile = session.get(); @@ -110,13 +116,6 @@ public class ParseSyslog extends AbstractProcessor { return; } - final String charsetName = context.getProperty(CHARSET).getValue(); - - // If the parser already exists and uses the same charset, it does not need to be re-initialized - if (parser == null || !parser.getCharsetName().equals(charsetName)) { - parser = new SyslogParser(Charset.forName(charsetName)); - } - final byte[] buffer = new byte[(int) flowFile.getSize()]; session.read(flowFile, in -> StreamUtils.fillBuffer(in, buffer)); diff --git a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/test/java/org/apache/nifi/processors/standard/TestParseSyslog.java b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/test/java/org/apache/nifi/processors/standard/TestParseSyslog.java index ad22454f435..94edd19728e 100644 --- a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/test/java/org/apache/nifi/processors/standard/TestParseSyslog.java +++ b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/test/java/org/apache/nifi/processors/standard/TestParseSyslog.java @@ -94,4 +94,25 @@ public class TestParseSyslog { runner.assertAllFlowFilesTransferred(ParseSyslog.REL_FAILURE, 1); } + + @Test + public void testConcurrentTasksShareScheduledParser() { + final TestRunner runner = TestRunners.newTestRunner(new ParseSyslog()); + final int numThreads = 8; + final int flowFiles = 32; + runner.setThreadCount(numThreads); + + for (int i = 0; i < flowFiles; i++) { + runner.enqueue(VALID_MESSAGE_RFC3164_0.getBytes()); + } + + runner.run(flowFiles); + + runner.assertAllFlowFilesTransferred(ParseSyslog.REL_SUCCESS, flowFiles); + for (final MockFlowFile mff : runner.getFlowFilesForRelationship(ParseSyslog.REL_SUCCESS)) { + mff.assertAttributeEquals(SyslogAttributes.SYSLOG_BODY.key(), BODY); + mff.assertAttributeEquals(SyslogAttributes.SYSLOG_HOSTNAME.key(), HOST); + mff.assertAttributeEquals(SyslogAttributes.SYSLOG_PRIORITY.key(), PRI); + } + } }
