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;
+        }
+    }
 }

Reply via email to