This is an automated email from the ASF dual-hosted git repository.
davsclaus pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/camel.git
The following commit(s) were added to refs/heads/main by this push:
new 323e370e67d1 CAMEL-24933, CAMEL-24934, CAMEL-16073: camel-pulsar -
consumer acknowledgement and exchange handling
323e370e67d1 is described below
commit 323e370e67d12eca1142bf5d5458a0e85c2700df
Author: Andrea Cosentino <[email protected]>
AuthorDate: Mon Sep 28 18:45:37 2026 +0200
CAMEL-24933, CAMEL-24934, CAMEL-16073: camel-pulsar - consumer
acknowledgement and exchange handling
CAMEL-24933: PulsarMessageUtils.updateExchange returned a copy of the
exchange, so the pooled exchange the consumer created was never released,
one per message. It now populates and returns the exchange it was given.
CAMEL-24934: a failed acknowledgement is now reported with its cause.
CAMEL-16073: when the route fails and allowManualAcknowledgement is false
(the default), the consumer negatively acknowledges the message instead of
leaving it to the acknowledgement timeout. It passes the Message, so a
configured negativeAckRedeliveryBackoff sees the redelivery count and
escalates. The nack is sent before the exception handler runs. Redelivery
now follows negativeAckRedeliveryDelayMicros (60 seconds by default); the
upgrade guide shows how to keep the previous 10 second timing.
Closes #26779
Co-authored-by: Claude Opus 5 <[email protected]>
---
.../component/pulsar/PulsarMessageListener.java | 21 ++-
.../pulsar/utils/message/PulsarMessageUtils.java | 15 +-
.../PulsarMessageListenerAcknowledgementTest.java | 197 +++++++++++++++++++++
.../utils/message/PulsarMessageUtilsTest.java | 34 ++++
.../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 24 +++
5 files changed, 282 insertions(+), 9 deletions(-)
diff --git
a/components/camel-pulsar/src/main/java/org/apache/camel/component/pulsar/PulsarMessageListener.java
b/components/camel-pulsar/src/main/java/org/apache/camel/component/pulsar/PulsarMessageListener.java
index 9038c354d06d..a3481bf770e7 100644
---
a/components/camel-pulsar/src/main/java/org/apache/camel/component/pulsar/PulsarMessageListener.java
+++
b/components/camel-pulsar/src/main/java/org/apache/camel/component/pulsar/PulsarMessageListener.java
@@ -50,14 +50,16 @@ public class PulsarMessageListener implements
MessageListener<byte[]> {
pulsarConsumer.getAsyncProcessor().process(exchange, doneSync -> {
try {
if (exchange.getException() != null) {
+ // tell the broker first: a custom ExceptionHandler that
throws must not cost us the
+ // negative acknowledgement, which would silently fall
back to ack-timeout redelivery
+ negativeAcknowledge(consumer, message);
pulsarConsumer.getExceptionHandler().handleException("Error processing
exchange", exchange,
exchange.getException());
} else {
try {
acknowledge(consumer, message);
} catch (Exception e) {
-
pulsarConsumer.getExceptionHandler().handleException("Error processing
exchange", exchange,
- exchange.getException());
+
pulsarConsumer.getExceptionHandler().handleException("Error acknowledging
message", exchange, e);
}
}
} finally {
@@ -73,4 +75,19 @@ public class PulsarMessageListener implements
MessageListener<byte[]> {
}
}
+ /**
+ * Tells the broker the message was not processed, so that it is
redelivered after
+ * {@code negativeAckRedeliveryDelayMicros} instead of waiting for the
acknowledgement timeout. Left to the route
+ * when manual acknowledgement is enabled, the same way {@link
#acknowledge} is.
+ * <p>
+ * The message is passed rather than its id on purpose: only that overload
carries the redelivery count into
+ * {@code NegativeAcksTracker}, and without it a configured {@code
negativeAckRedeliveryBackoff} is always asked for
+ * the delay of attempt zero and never escalates.
+ */
+ private void negativeAcknowledge(final Consumer<byte[]> consumer, final
Message<byte[]> message) {
+ if (!endpoint.getPulsarConfiguration().isAllowManualAcknowledgement())
{
+ consumer.negativeAcknowledge(message);
+ }
+ }
+
}
diff --git
a/components/camel-pulsar/src/main/java/org/apache/camel/component/pulsar/utils/message/PulsarMessageUtils.java
b/components/camel-pulsar/src/main/java/org/apache/camel/component/pulsar/utils/message/PulsarMessageUtils.java
index d03b3fe1287c..a9ff18556622 100644
---
a/components/camel-pulsar/src/main/java/org/apache/camel/component/pulsar/utils/message/PulsarMessageUtils.java
+++
b/components/camel-pulsar/src/main/java/org/apache/camel/component/pulsar/utils/message/PulsarMessageUtils.java
@@ -41,10 +41,13 @@ public final class PulsarMessageUtils {
private PulsarMessageUtils() {
}
- public static Exchange updateExchange(final Message<byte[]> message, final
Exchange input) {
- final Exchange output = input.copy();
-
- org.apache.camel.Message msg = output.getIn();
+ /**
+ * Populates the given exchange from the Pulsar message and returns it.
The exchange is updated in place: copying it
+ * would orphan the exchange the consumer took from the exchange factory,
which with a pooled factory is never
+ * returned to the pool.
+ */
+ public static Exchange updateExchange(final Message<byte[]> message, final
Exchange exchange) {
+ org.apache.camel.Message msg = exchange.getIn();
msg.setHeader(EVENT_TIME, message.getEventTime());
msg.setHeader(MESSAGE_ID, message.getMessageId());
@@ -60,9 +63,7 @@ public final class PulsarMessageUtils {
msg.setBody(message.getValue());
- output.setIn(msg);
-
- return output;
+ return exchange;
}
public static Exchange updateExchangeWithException(final Exception
exception, final Exchange input) {
diff --git
a/components/camel-pulsar/src/test/java/org/apache/camel/component/pulsar/PulsarMessageListenerAcknowledgementTest.java
b/components/camel-pulsar/src/test/java/org/apache/camel/component/pulsar/PulsarMessageListenerAcknowledgementTest.java
new file mode 100644
index 000000000000..85261d991b10
--- /dev/null
+++
b/components/camel-pulsar/src/test/java/org/apache/camel/component/pulsar/PulsarMessageListenerAcknowledgementTest.java
@@ -0,0 +1,197 @@
+/*
+ * 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.nio.charset.StandardCharsets;
+import java.util.Collections;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+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.Message;
+import org.apache.pulsar.client.api.MessageId;
+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.doThrow;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.timeout;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/**
+ * A route failure used to leave the message neither acknowledged nor
negatively acknowledged, so it was only
+ * redelivered once the acknowledgement timeout expired (CAMEL-16073), and a
failing acknowledgement was reported
+ * without its cause.
+ */
+public class PulsarMessageListenerAcknowledgementTest extends CamelTestSupport
{
+
+ private static final String ENDPOINT_URI
+ =
"pulsar:persistent://public/default/camel-ack-test?subscriptionName=camel-subscription";
+
+ private final MessageId messageId = mock(MessageId.class);
+
+ @Test
+ public void testSuccessfulExchangeIsAcknowledged() throws Exception {
+ final Consumer<byte[]> pulsarConsumer = mock(Consumer.class);
+
+ listener().received(pulsarConsumer, message("ok"));
+
+ verify(pulsarConsumer, timeout(5000)).acknowledge(messageId);
+ verify(pulsarConsumer, never()).negativeAcknowledge(messageId);
+ }
+
+ @Test
+ public void testFailedExchangeIsNegativelyAcknowledged() throws Exception {
+ final Consumer<byte[]> pulsarConsumer = mock(Consumer.class);
+ final Message<byte[]> message = message("fail");
+
+ listener().received(pulsarConsumer, message);
+
+ // the Message overload, not the MessageId one: only that carries the
redelivery count into
+ // NegativeAcksTracker, so a configured negativeAckRedeliveryBackoff
can escalate
+ verify(pulsarConsumer, timeout(5000)).negativeAcknowledge(message);
+ verify(pulsarConsumer, never()).negativeAcknowledge(messageId);
+ verify(pulsarConsumer, never()).acknowledge(messageId);
+ }
+
+ @Test
+ public void testManualAcknowledgementLeavesTheFailedMessageToTheRoute()
throws Exception {
+ final Consumer<byte[]> pulsarConsumer = mock(Consumer.class);
+ final Message<byte[]> message = message("fail");
+
+ final CapturingExceptionHandler exceptionHandler = new
CapturingExceptionHandler();
+ pulsarConsumer().setExceptionHandler(exceptionHandler);
+
+ context.getEndpoint(ENDPOINT_URI,
PulsarEndpoint.class).getPulsarConfiguration()
+ .setAllowManualAcknowledgement(true);
+ try {
+ listener().received(pulsarConsumer, message);
+
+ // wait for the callback to have run before asserting that nothing
was sent to the broker
+ assertTrue(exceptionHandler.latch.await(5, TimeUnit.SECONDS),
+ "the route failure should still be reported");
+ verify(pulsarConsumer, never()).negativeAcknowledge(message);
+ verify(pulsarConsumer, never()).negativeAcknowledge(messageId);
+ verify(pulsarConsumer, never()).acknowledge(messageId);
+ } finally {
+ context.getEndpoint(ENDPOINT_URI,
PulsarEndpoint.class).getPulsarConfiguration()
+ .setAllowManualAcknowledgement(false);
+ }
+ }
+
+ @Test
+ public void testAcknowledgeFailureIsReportedWithItsCause() throws
Exception {
+ final Consumer<byte[]> pulsarConsumer = mock(Consumer.class);
+ final PulsarClientException failure = new
PulsarClientException("cannot acknowledge");
+ doThrow(failure).when(pulsarConsumer).acknowledge(messageId);
+
+ final CapturingExceptionHandler exceptionHandler = new
CapturingExceptionHandler();
+ pulsarConsumer().setExceptionHandler(exceptionHandler);
+
+ listener().received(pulsarConsumer, message("ok"));
+
+ assertTrue(exceptionHandler.latch.await(5, TimeUnit.SECONDS), "the
acknowledgement failure should be reported");
+ assertNotNull(exceptionHandler.captured.get(), "the acknowledgement
failure should carry its cause");
+ assertEquals("cannot acknowledge",
exceptionHandler.captured.get().getMessage());
+ }
+
+ private PulsarMessageListener listener() {
+ return new PulsarMessageListener(context.getEndpoint(ENDPOINT_URI,
PulsarEndpoint.class), pulsarConsumer());
+ }
+
+ private PulsarConsumer pulsarConsumer() {
+ return (PulsarConsumer) context.getRoutes().get(0).getConsumer();
+ }
+
+ private Message<byte[]> message(String body) {
+ final Message<byte[]> message = mock(Message.class);
+
when(message.getValue()).thenReturn(body.getBytes(StandardCharsets.UTF_8));
+ when(message.getMessageId()).thenReturn(messageId);
+ when(message.getProperties()).thenReturn(Collections.emptyMap());
+ return message;
+ }
+
+ @Override
+ protected CamelContext createCamelContext() throws Exception {
+ final CamelContext context = super.createCamelContext();
+
+ final ConsumerBuilder<byte[]> builder = mock(ConsumerBuilder.class,
RETURNS_SELF);
+ when(builder.subscribe()).thenReturn(mock(Consumer.class));
+
+ 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() {
+ from(ENDPOINT_URI)
+ .process(exchange -> {
+ if
("fail".equals(exchange.getIn().getBody(String.class))) {
+ throw new IllegalStateException("simulated
route failure");
+ }
+ })
+ .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();
+ }
+ }
+}
diff --git
a/components/camel-pulsar/src/test/java/org/apache/camel/component/pulsar/utils/message/PulsarMessageUtilsTest.java
b/components/camel-pulsar/src/test/java/org/apache/camel/component/pulsar/utils/message/PulsarMessageUtilsTest.java
index 18deba91db0f..1e98017524ce 100644
---
a/components/camel-pulsar/src/test/java/org/apache/camel/component/pulsar/utils/message/PulsarMessageUtilsTest.java
+++
b/components/camel-pulsar/src/test/java/org/apache/camel/component/pulsar/utils/message/PulsarMessageUtilsTest.java
@@ -17,10 +17,20 @@
package org.apache.camel.component.pulsar.utils.message;
import java.io.Serializable;
+import java.util.Collections;
+import org.apache.camel.Exchange;
+import org.apache.camel.impl.DefaultCamelContext;
+import org.apache.camel.support.DefaultExchange;
+import org.apache.pulsar.client.api.Message;
+import org.apache.pulsar.client.api.MessageId;
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.assertSame;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
public class PulsarMessageUtilsTest {
@@ -48,6 +58,30 @@ public class PulsarMessageUtilsTest {
assertNotNull(expected);
}
+
+ @Test
+ public void testUpdateExchangeUpdatesTheGivenExchangeInPlace() {
+ final DefaultCamelContext context = new DefaultCamelContext();
+ final Exchange exchange = new DefaultExchange(context);
+
+ final MessageId messageId = mock(MessageId.class);
+ final Message<byte[]> message = mock(Message.class);
+ when(message.getValue()).thenReturn("Hello World!".getBytes());
+ when(message.getMessageId()).thenReturn(messageId);
+ when(message.getKey()).thenReturn("aKey");
+ when(message.getTopicName()).thenReturn("aTopic");
+ when(message.getProducerName()).thenReturn("aProducer");
+ when(message.getProperties()).thenReturn(Collections.emptyMap());
+
+ final Exchange updated = PulsarMessageUtils.updateExchange(message,
exchange);
+
+ // a copy here would orphan the exchange the consumer took from the
exchange factory
+ assertSame(exchange, updated);
+ assertEquals(messageId,
updated.getIn().getHeader(PulsarMessageHeaders.MESSAGE_ID));
+ assertEquals("aKey",
updated.getIn().getHeader(PulsarMessageHeaders.KEY));
+ assertEquals("aTopic",
updated.getIn().getHeader(PulsarMessageHeaders.TOPIC_NAME));
+ assertEquals("Hello World!", updated.getIn().getBody(String.class));
+ }
}
class Obj implements Serializable {
diff --git
a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
index 7afd41e043c9..22e3d3ae84a2 100644
--- a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
+++ b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
@@ -3192,6 +3192,30 @@ reply of a discarded exchange (`waitForTaskToComplete`)
is released with that ex
`timeout`, or forever when the timeout is disabled. On completions handed over
to a discarded InOnly exchange, such as
the commit or rollback of the consumer that received the message, now run as a
failure, where previously they never ran.
+=== camel-pulsar - PulsarMessageUtils.updateExchange returns the exchange it
was given
+
+`PulsarMessageUtils.updateExchange(message, exchange)` used to return a *copy*
of the exchange passed to
+it. It now populates and returns that same instance.
+
+The copy orphaned the exchange the consumer had taken from the exchange
factory, so with
+`camel.main.exchange-factory=pooled` every consumed message leaked one pooled
exchange and the pool never
+refilled. Code outside the component that called this method and relied on
getting an independent copy
+must make its own copy instead.
+
+=== camel-pulsar - a failed exchange is negatively acknowledged
+
+When a route fails, the consumer now calls `negativeAcknowledge` on the Pulsar
consumer instead of
+leaving the message unacknowledged. This only applies when
`allowManualAcknowledgement` is `false`
+(the default); with manual acknowledgement the route stays in charge, as
before.
+
+This changes when the message comes back. Previously it was redelivered once
the acknowledgement
+timeout expired, which `camel-pulsar` sets to 10 seconds by default through
`ackTimeoutMillis`. A
+negative acknowledgement removes the message from the client's
unacknowledged-message tracker, so
+redelivery now follows `negativeAckRedeliveryDelayMicros`, which defaults to
60 seconds, and honours
+`negativeAckRedeliveryBackoff` when one is configured.
+
+A route that wants the previous timing can set
`negativeAckRedeliveryDelayMicros=10000000`.
+
=== camel-seda - multipleConsumers broadcasts to consumers with different uri
options
With `multipleConsumers=true` every consumer of a SEDA queue now receives a
copy of each message, also when the consumers