This is an automated email from the ASF dual-hosted git repository.

gnodet pushed a commit to branch camel-4.22.x
in repository https://gitbox.apache.org/repos/asf/camel.git


The following commit(s) were added to refs/heads/camel-4.22.x by this push:
     new 097b69b92929 [backport camel-4.22.x] CAMEL-24932: camel-pulsar - do 
not kill the polling loop on a receive error (#26887)
097b69b92929 is described below

commit 097b69b929295bb34c2409b3d5d6f1019b611aaa
Author: Guillaume Nodet <[email protected]>
AuthorDate: Fri Sep 25 11:33:40 2026 +0200

    [backport camel-4.22.x] CAMEL-24932: camel-pulsar - do not kill the polling 
loop on a receive error (#26887)
---
 .../camel/component/pulsar/PulsarConsumer.java     |  33 ++++-
 .../pulsar/PulsarConsumerReceiveErrorTest.java     | 142 +++++++++++++++++++++
 2 files changed, 173 insertions(+), 2 deletions(-)

diff --git 
a/components/camel-pulsar/src/main/java/org/apache/camel/component/pulsar/PulsarConsumer.java
 
b/components/camel-pulsar/src/main/java/org/apache/camel/component/pulsar/PulsarConsumer.java
index 62542de2f674..b4890387c9b8 100644
--- 
a/components/camel-pulsar/src/main/java/org/apache/camel/component/pulsar/PulsarConsumer.java
+++ 
b/components/camel-pulsar/src/main/java/org/apache/camel/component/pulsar/PulsarConsumer.java
@@ -41,6 +41,11 @@ import static 
org.apache.camel.component.pulsar.utils.PulsarUtils.stopExecutors;
 public class PulsarConsumer extends DefaultConsumer implements Suspendable {
     private static final Logger LOGGER = 
LoggerFactory.getLogger(PulsarConsumer.class);
 
+    /**
+     * How long a consumer thread waits after a failure it cannot classify, 
before polling again.
+     */
+    private static final long ERROR_RETRY_DELAY_MILLIS = 1000;
+
     private final PulsarEndpoint pulsarEndpoint;
     private final ConsumerCreationStrategyFactory 
consumerCreationStrategyFactory;
 
@@ -128,6 +133,10 @@ public class PulsarConsumer extends DefaultConsumer 
implements Suspendable {
                 try {
                     Message<byte[]> msg = consumer.receive();
                     listener.received(consumer, msg);
+                } catch (PulsarClientException.AlreadyClosedException e) {
+                    // the consumer is gone, so there is nothing left for this 
loop to poll
+                    LOGGER.info("Pulsar consumer is closed, exiting");
+                    running = false;
                 } catch (PulsarClientException e) {
                     if (e.getCause() instanceof InterruptedException) {
                         // this means that our executor is shutting down
@@ -136,12 +145,32 @@ public class PulsarConsumer extends DefaultConsumer 
implements Suspendable {
                         // by exiting the loop. We make it explicit instead of 
breaking the loop.
                         running = false;
                     } else {
-                        endpoint.getExceptionHandler().handleException(e);
+                        getExceptionHandler().handleException("Error consuming 
from pulsar", e);
+                        running = waitBeforeRetry();
                     }
                 } catch (Exception e) {
-                    endpoint.getExceptionHandler().handleException(e);
+                    getExceptionHandler().handleException("Error consuming 
from pulsar", e);
+                    running = waitBeforeRetry();
                 }
             }
         }
+
+        /**
+         * Waits before polling again, so that a failure which does not clear 
- an unreachable broker, a consumer in
+         * Failed state - does not turn this into a hot loop on every consumer 
thread, reporting the same error to the
+         * exception handler as fast as the CPU allows.
+         *
+         * @return {@code false} when the wait was interrupted, which is how a 
stopping consumer leaves the loop
+         */
+        private boolean waitBeforeRetry() {
+            try {
+                Thread.sleep(ERROR_RETRY_DELAY_MILLIS);
+                return true;
+            } catch (InterruptedException e) {
+                Thread.currentThread().interrupt();
+                LOGGER.info("Received shutdown signal, exiting");
+                return false;
+            }
+        }
     }
 }
diff --git 
a/components/camel-pulsar/src/test/java/org/apache/camel/component/pulsar/PulsarConsumerReceiveErrorTest.java
 
b/components/camel-pulsar/src/test/java/org/apache/camel/component/pulsar/PulsarConsumerReceiveErrorTest.java
new file mode 100644
index 000000000000..db2bfceef9b6
--- /dev/null
+++ 
b/components/camel-pulsar/src/test/java/org/apache/camel/component/pulsar/PulsarConsumerReceiveErrorTest.java
@@ -0,0 +1,142 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *      http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.camel.component.pulsar;
+
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicReference;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.Exchange;
+import org.apache.camel.RoutesBuilder;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.spi.ExceptionHandler;
+import org.apache.camel.test.junit6.CamelTestSupport;
+import org.apache.pulsar.client.api.Consumer;
+import org.apache.pulsar.client.api.ConsumerBuilder;
+import org.apache.pulsar.client.api.PulsarClient;
+import org.apache.pulsar.client.api.PulsarClientException;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Answers.RETURNS_SELF;
+import static org.mockito.Mockito.after;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.timeout;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/**
+ * A receive error used to be reported through the endpoint exception handler, 
which is null unless the route configures
+ * one, so the polling loop died on a NullPointerException and that consumer 
stopped consuming in silence.
+ */
+public class PulsarConsumerReceiveErrorTest extends CamelTestSupport {
+
+    private static final String ROUTE_ID = "pulsar-receive-error";
+
+    private Consumer<byte[]> pulsarConsumer;
+
+    private final AtomicInteger receiveCalls = new AtomicInteger();
+
+    @Test
+    public void testReceiveErrorIsHandledAndTheLoopSurvives() throws Exception 
{
+        final PulsarConsumer consumer = (PulsarConsumer) 
context.getRoute(ROUTE_ID).getConsumer();
+
+        final CapturingExceptionHandler exceptionHandler = new 
CapturingExceptionHandler();
+        consumer.setExceptionHandler(exceptionHandler);
+
+        context.getRouteController().startRoute(ROUTE_ID);
+
+        assertTrue(exceptionHandler.latch.await(10, TimeUnit.SECONDS),
+                "the receive error should be handed to the consumer exception 
handler");
+        assertNotNull(exceptionHandler.captured.get(), "the failure cause 
should be reported");
+        assertEquals("simulated receive failure", 
exceptionHandler.captured.get().getMessage());
+
+        // the loop must still be polling: the first call failed, so a second 
one proves it did not die
+        verify(pulsarConsumer, timeout(10000).times(2)).receive();
+
+        // the second call reported a closed consumer, which is terminal: the 
loop must leave rather than
+        // spin on an error that will never clear
+        verify(pulsarConsumer, after(1500).times(2)).receive();
+    }
+
+    @Override
+    protected CamelContext createCamelContext() throws Exception {
+        final CamelContext context = super.createCamelContext();
+
+        pulsarConsumer = mock(Consumer.class);
+        when(pulsarConsumer.receive()).thenAnswer(invocation -> {
+            if (receiveCalls.incrementAndGet() == 1) {
+                throw new PulsarClientException("simulated receive failure");
+            }
+            // a failure that never clears, which the loop must treat as 
terminal
+            throw new PulsarClientException.AlreadyClosedException("consumer 
is closed");
+        });
+
+        final ConsumerBuilder<byte[]> builder = mock(ConsumerBuilder.class, 
RETURNS_SELF);
+        when(builder.subscribe()).thenReturn(pulsarConsumer);
+
+        final PulsarClient pulsarClient = mock(PulsarClient.class);
+        when(pulsarClient.newConsumer()).thenReturn(builder);
+
+        final PulsarComponent component = new PulsarComponent(context);
+        component.setPulsarClient(pulsarClient);
+        context.addComponent("pulsar", component);
+
+        return context;
+    }
+
+    @Override
+    protected RoutesBuilder createRouteBuilder() {
+        return new RouteBuilder() {
+            @Override
+            public void configure() {
+                // started by the test, so that the exception handler is in 
place before the loop runs;
+                // no exceptionHandler option on the endpoint, which is 
exactly the case that used to NPE
+                from("pulsar:persistent://public/default/camel-receive-error"
+                     + 
"?messageListener=false&subscriptionName=camel-subscription")
+                        .routeId(ROUTE_ID).autoStartup(false)
+                        .to("mock:result");
+            }
+        };
+    }
+
+    private static final class CapturingExceptionHandler implements 
ExceptionHandler {
+
+        private final CountDownLatch latch = new CountDownLatch(1);
+        private final AtomicReference<Throwable> captured = new 
AtomicReference<>();
+
+        @Override
+        public void handleException(Throwable exception) {
+            handleException(null, null, exception);
+        }
+
+        @Override
+        public void handleException(String message, Throwable exception) {
+            handleException(message, null, exception);
+        }
+
+        @Override
+        public void handleException(String message, Exchange exchange, 
Throwable exception) {
+            captured.compareAndSet(null, exception);
+            latch.countDown();
+        }
+    }
+}

Reply via email to