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