This is an automated email from the ASF dual-hosted git repository.
turcsanyip pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/nifi.git
The following commit(s) were added to refs/heads/main by this push:
new cd2a3461404 NIFI-15935 Clean up AMQP resources after processor
failures (#11243)
cd2a3461404 is described below
commit cd2a346140481a7d4c0d8076c5bef70b20e7c7d9
Author: ing-mattioni <[email protected]>
AuthorDate: Mon Jul 20 18:47:36 2026 +0200
NIFI-15935 Clean up AMQP resources after processor failures (#11243)
---
.../amqp/processors/AbstractAMQPProcessor.java | 17 ++++++--
.../amqp/processors/AbstractAMQPProcessorTest.java | 50 ++++++++++++++++++++++
2 files changed, 64 insertions(+), 3 deletions(-)
diff --git
a/nifi-extension-bundles/nifi-amqp-bundle/nifi-amqp-processors/src/main/java/org/apache/nifi/amqp/processors/AbstractAMQPProcessor.java
b/nifi-extension-bundles/nifi-amqp-bundle/nifi-amqp-processors/src/main/java/org/apache/nifi/amqp/processors/AbstractAMQPProcessor.java
index f8f6b8f8a58..0aa7aba153e 100644
---
a/nifi-extension-bundles/nifi-amqp-bundle/nifi-amqp-processors/src/main/java/org/apache/nifi/amqp/processors/AbstractAMQPProcessor.java
+++
b/nifi-extension-bundles/nifi-amqp-bundle/nifi-amqp-processors/src/main/java/org/apache/nifi/amqp/processors/AbstractAMQPProcessor.java
@@ -244,8 +244,10 @@ abstract class AbstractAMQPProcessor<T extends AMQPWorker>
extends AbstractProce
context.yield();
closeResource(resource);
} catch (Exception e) {
- getLogger().error("Processor failure", e);
+ getLogger().error("Processor failure, dropping the client", e);
+ session.rollback();
context.yield();
+ closeResource(resource);
}
}
@@ -283,9 +285,10 @@ abstract class AbstractAMQPProcessor<T extends AMQPWorker>
extends AbstractProce
private AMQPResource<T> createResource(final ProcessContext context) {
Connection connection = null;
+ ExecutorService executor = null;
try {
- ExecutorService executor =
Executors.newSingleThreadExecutor(BasicThreadFactory.builder()
- .namingPattern("AMQP Consumer: " + getIdentifier())
+ executor =
Executors.newSingleThreadExecutor(BasicThreadFactory.builder()
+ .namingPattern("AMQP Client: " + getIdentifier() + "-%d")
.build());
connection = createConnection(context, executor);
T worker = createAMQPWorker(context, connection);
@@ -298,6 +301,14 @@ abstract class AbstractAMQPProcessor<T extends AMQPWorker>
extends AbstractProce
getLogger().error("Failed to close AMQP Connection",
closingEx);
}
}
+ if (executor != null) {
+ try {
+ executor.shutdown();
+ } catch (Exception closingEx) {
+ getLogger().error("Failed to shut down AMQP Executor",
closingEx);
+ e.addSuppressed(closingEx);
+ }
+ }
throw e;
}
}
diff --git
a/nifi-extension-bundles/nifi-amqp-bundle/nifi-amqp-processors/src/test/java/org/apache/nifi/amqp/processors/AbstractAMQPProcessorTest.java
b/nifi-extension-bundles/nifi-amqp-bundle/nifi-amqp-processors/src/test/java/org/apache/nifi/amqp/processors/AbstractAMQPProcessorTest.java
index 096c168eea5..be76f6e7828 100644
---
a/nifi-extension-bundles/nifi-amqp-bundle/nifi-amqp-processors/src/test/java/org/apache/nifi/amqp/processors/AbstractAMQPProcessorTest.java
+++
b/nifi-extension-bundles/nifi-amqp-bundle/nifi-amqp-processors/src/test/java/org/apache/nifi/amqp/processors/AbstractAMQPProcessorTest.java
@@ -16,6 +16,9 @@
*/
package org.apache.nifi.amqp.processors;
+import com.rabbitmq.client.Connection;
+import org.apache.nifi.processor.ProcessContext;
+import org.apache.nifi.processor.exception.ProcessException;
import org.apache.nifi.reporting.InitializationException;
import org.apache.nifi.ssl.SSLContextProvider;
import org.apache.nifi.util.TestRunner;
@@ -23,7 +26,11 @@ import org.apache.nifi.util.TestRunners;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
+import java.util.concurrent.ExecutorService;
+
+import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
/**
@@ -88,6 +95,21 @@ public class AbstractAMQPProcessorTest {
testRunner.assertNotValid();
}
+ @Test
+ public void testExecutorShutdownWhenResourceCreationFails() throws
Exception {
+ final FailingWorkerCreationConsumeAMQP processor = new
FailingWorkerCreationConsumeAMQP();
+ final TestRunner runner = TestRunners.newTestRunner(processor);
+ runner.setProperty(ConsumeAMQP.QUEUE, "queue");
+ runner.setProperty(AbstractAMQPProcessor.BROKERS, "localhost:5672");
+ runner.setProperty(AbstractAMQPProcessor.USER, "user");
+ runner.setProperty(AbstractAMQPProcessor.PASSWORD, "password");
+
+ runner.run();
+
+ assertTrue(processor.getExecutor().isShutdown());
+ verify(processor.getConnection()).close();
+ }
+
private void configureSSLContextService() throws InitializationException {
SSLContextProvider sslContextProvider = mock(SSLContextProvider.class);
when(sslContextProvider.getIdentifier()).thenReturn("ssl-context");
@@ -96,4 +118,32 @@ public class AbstractAMQPProcessorTest {
testRunner.setProperty(AbstractAMQPProcessor.SSL_CONTEXT_SERVICE,
"ssl-context");
}
+
+ private static class FailingWorkerCreationConsumeAMQP extends ConsumeAMQP {
+ private final Connection connection = mock(Connection.class);
+ private ExecutorService executor;
+
+ private FailingWorkerCreationConsumeAMQP() {
+ when(connection.isOpen()).thenReturn(true);
+ }
+
+ @Override
+ protected Connection createConnection(final ProcessContext context,
final ExecutorService executor) {
+ this.executor = executor;
+ return connection;
+ }
+
+ @Override
+ protected AMQPConsumer createAMQPWorker(final ProcessContext context,
final Connection connection) {
+ throw new ProcessException("Worker creation failed");
+ }
+
+ private Connection getConnection() {
+ return connection;
+ }
+
+ private ExecutorService getExecutor() {
+ return executor;
+ }
+ }
}