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 73cdd7aad837 CAMEL-25272: camel-rocketmq - consume a message again
when its route failed, and fail an InOut send without replyToTopic (#27292)
73cdd7aad837 is described below
commit 73cdd7aad83750bd68712dc6f96664e27c30cdb9
Author: allthingssecurity <[email protected]>
AuthorDate: Sat Oct 3 12:19:15 2026 +0530
CAMEL-25272: camel-rocketmq - consume a message again when its route
failed, and fail an InOut send without replyToTopic (#27292)
Two places where camel-rocketmq reports a failure as success:
- **Consumer.** `RocketMQConsumer`'s listener answered `CONSUME_SUCCESS`
whenever the processor returned, and `RECONSUME_LATER` only when it threw. A
failing route does not throw: the error handler sets the exception on the
exchange. So the message of every failed exchange was acknowledged and lost (no
retry topic, no dead letter queue).
Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
.../camel/catalog/docs/rocketmq-component.adoc | 8 ++
.../src/main/docs/rocketmq-component.adoc | 8 ++
.../camel/component/rocketmq/RocketMQConsumer.java | 36 ++++---
.../camel/component/rocketmq/RocketMQProducer.java | 5 +-
.../rocketmq/RocketMQConsumerFailureTest.java | 104 +++++++++++++++++++++
.../rocketmq/RocketMQProducerSendFailureTest.java | 48 ++++++++++
.../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 9 ++
7 files changed, 205 insertions(+), 13 deletions(-)
diff --git
a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/rocketmq-component.adoc
b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/rocketmq-component.adoc
index 41d6ecacf0a4..0d76ac38adab 100644
---
a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/rocketmq-component.adoc
+++
b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/rocketmq-component.adoc
@@ -79,6 +79,14 @@ from("rocketmq:START_TOPIC?consumerGroup=c1")
.to("log:InOutRoute?showAll=true")
----
+=== Error handling
+
+The consumer acknowledges a message to RocketMQ (`CONSUME_SUCCESS`) when the
exchange completes. When the exchange
+fails, or is marked rollback only, the consumer answers `RECONSUME_LATER`, and
RocketMQ delivers the message again with
+an increasing delay, up to the maximum reconsume times of the consumer group
(16 by default); after that the message
+goes to the dead letter queue of the consumer group. To acknowledge a message
whose processing failed, handle the
+exception in the route, for example with `onException(...).handled(true)`.
+
== Examples
Receive messages from a topic named `from_topic`, route to `to_topic`.
diff --git a/components/camel-rocketmq/src/main/docs/rocketmq-component.adoc
b/components/camel-rocketmq/src/main/docs/rocketmq-component.adoc
index 41d6ecacf0a4..0d76ac38adab 100644
--- a/components/camel-rocketmq/src/main/docs/rocketmq-component.adoc
+++ b/components/camel-rocketmq/src/main/docs/rocketmq-component.adoc
@@ -79,6 +79,14 @@ from("rocketmq:START_TOPIC?consumerGroup=c1")
.to("log:InOutRoute?showAll=true")
----
+=== Error handling
+
+The consumer acknowledges a message to RocketMQ (`CONSUME_SUCCESS`) when the
exchange completes. When the exchange
+fails, or is marked rollback only, the consumer answers `RECONSUME_LATER`, and
RocketMQ delivers the message again with
+an increasing delay, up to the maximum reconsume times of the consumer group
(16 by default); after that the message
+goes to the dead letter queue of the consumer group. To acknowledge a message
whose processing failed, handle the
+exception in the route, for example with `onException(...).handled(true)`.
+
== Examples
Receive messages from a topic named `from_topic`, route to `to_topic`.
diff --git
a/components/camel-rocketmq/src/main/java/org/apache/camel/component/rocketmq/RocketMQConsumer.java
b/components/camel-rocketmq/src/main/java/org/apache/camel/component/rocketmq/RocketMQConsumer.java
index c26dd2e10cb1..e00deba00ecc 100644
---
a/components/camel-rocketmq/src/main/java/org/apache/camel/component/rocketmq/RocketMQConsumer.java
+++
b/components/camel-rocketmq/src/main/java/org/apache/camel/component/rocketmq/RocketMQConsumer.java
@@ -17,6 +17,8 @@
package org.apache.camel.component.rocketmq;
+import java.util.List;
+
import org.apache.camel.Exchange;
import org.apache.camel.Processor;
import org.apache.camel.Suspendable;
@@ -24,6 +26,7 @@ import org.apache.camel.support.DefaultConsumer;
import org.apache.rocketmq.client.AccessChannel;
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.MessageSelector;
+import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import
org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import org.apache.rocketmq.client.exception.MQClientException;
@@ -62,21 +65,30 @@ public class RocketMQConsumer extends DefaultConsumer
implements Suspendable {
mqPushConsumer.setAccessChannel(AccessChannel.valueOf(endpoint.getAccessChannel()));
mqPushConsumer.subscribe(endpoint.getTopicName(), messageSelector);
- mqPushConsumer.registerMessageListener((MessageListenerConcurrently)
(msgs, context) -> {
- MessageExt messageExt = msgs.get(0);
- Exchange exchange =
endpoint.createRocketExchange(messageExt.getBody());
-
RocketMQMessageConverter.populateHeadersByMessageExt(exchange.getIn(),
messageExt);
- try {
- getProcessor().process(exchange);
- } catch (Exception e) {
- getExceptionHandler().handleException(e);
- return ConsumeConcurrentlyStatus.RECONSUME_LATER;
- }
- return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
- });
+ mqPushConsumer.registerMessageListener((MessageListenerConcurrently)
this::consumeMessage);
mqPushConsumer.start();
}
+ ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs,
ConsumeConcurrentlyContext context) {
+ MessageExt messageExt = msgs.get(0);
+ Exchange exchange =
endpoint.createRocketExchange(messageExt.getBody());
+ RocketMQMessageConverter.populateHeadersByMessageExt(exchange.getIn(),
messageExt);
+ try {
+ getProcessor().process(exchange);
+ } catch (Exception e) {
+ exchange.setException(e);
+ }
+ // a failed route sets the exception on the exchange (or marks it
rollback only): the message must then be
+ // consumed again, acknowledging it would lose it
+ if (exchange.getException() != null || exchange.isRollbackOnly() ||
exchange.isRollbackOnlyLast()) {
+ if (exchange.getException() != null) {
+ getExceptionHandler().handleException("Error processing
exchange", exchange, exchange.getException());
+ }
+ return ConsumeConcurrentlyStatus.RECONSUME_LATER;
+ }
+ return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
+ }
+
private void stopConsumer() {
if (mqPushConsumer != null) {
mqPushConsumer.shutdown();
diff --git
a/components/camel-rocketmq/src/main/java/org/apache/camel/component/rocketmq/RocketMQProducer.java
b/components/camel-rocketmq/src/main/java/org/apache/camel/component/rocketmq/RocketMQProducer.java
index e6a181b2cc79..b825647ba207 100644
---
a/components/camel-rocketmq/src/main/java/org/apache/camel/component/rocketmq/RocketMQProducer.java
+++
b/components/camel-rocketmq/src/main/java/org/apache/camel/component/rocketmq/RocketMQProducer.java
@@ -125,8 +125,11 @@ public class RocketMQProducer extends DefaultAsyncProducer
{
@Override
public void onException(Throwable e) {
try {
- replyManager.cancelMessageKey(generateKey);
exchange.setException(e);
+ // there is no reply manager when replyToTopic is not set
+ if (replyManager != null) {
+ replyManager.cancelMessageKey(generateKey);
+ }
} finally {
callback.done(false);
}
diff --git
a/components/camel-rocketmq/src/test/java/org/apache/camel/component/rocketmq/RocketMQConsumerFailureTest.java
b/components/camel-rocketmq/src/test/java/org/apache/camel/component/rocketmq/RocketMQConsumerFailureTest.java
new file mode 100644
index 000000000000..f132b4533ef7
--- /dev/null
+++
b/components/camel-rocketmq/src/test/java/org/apache/camel/component/rocketmq/RocketMQConsumerFailureTest.java
@@ -0,0 +1,104 @@
+/*
+ * 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.rocketmq;
+
+import java.nio.charset.StandardCharsets;
+import java.util.List;
+
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.apache.camel.test.junit6.CamelTestSupport;
+import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
+import org.apache.rocketmq.common.message.MessageExt;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
+/**
+ * The status returned to the RocketMQ push consumer acknowledges the message:
a message whose route failed must be
+ * consumed again, not acknowledged.
+ */
+public class RocketMQConsumerFailureTest extends CamelTestSupport {
+
+ private static final String ROCKETMQ_URI =
"rocketmq:START_TOPIC?namesrvAddr=localhost:9876&consumerGroup=c1";
+
+ @Test
+ public void testFailedExchangeIsConsumedLater() throws Exception {
+ RocketMQConsumer consumer = createConsumer("direct:fail");
+
+ assertEquals(ConsumeConcurrentlyStatus.RECONSUME_LATER,
consumer.consumeMessage(List.of(message()), null));
+ }
+
+ @Test
+ public void testExceptionFromProcessorIsConsumedLater() throws Exception {
+ RocketMQEndpoint endpoint = context.getEndpoint(ROCKETMQ_URI,
RocketMQEndpoint.class);
+ RocketMQConsumer consumer = (RocketMQConsumer)
endpoint.createConsumer(exchange -> {
+ throw new IllegalStateException("Forced");
+ });
+
+ assertEquals(ConsumeConcurrentlyStatus.RECONSUME_LATER,
consumer.consumeMessage(List.of(message()), null));
+ }
+
+ @Test
+ public void testRollbackOnlyExchangeIsConsumedLater() throws Exception {
+ RocketMQConsumer consumer = createConsumer("direct:rollback");
+
+ assertEquals(ConsumeConcurrentlyStatus.RECONSUME_LATER,
consumer.consumeMessage(List.of(message()), null));
+ }
+
+ @Test
+ public void testCompletedExchangeIsAcknowledged() throws Exception {
+ MockEndpoint result = getMockEndpoint("mock:result");
+ result.expectedBodiesReceived("Hello");
+
+ RocketMQConsumer consumer = createConsumer("direct:ok");
+
+ assertEquals(ConsumeConcurrentlyStatus.CONSUME_SUCCESS,
consumer.consumeMessage(List.of(message()), null));
+ result.assertIsSatisfied();
+ }
+
+ private RocketMQConsumer createConsumer(String route) throws Exception {
+ RocketMQEndpoint endpoint = context.getEndpoint(ROCKETMQ_URI,
RocketMQEndpoint.class);
+ // the route sets the exception on the exchange when it fails, as the
route of a consumer does
+ return (RocketMQConsumer) endpoint.createConsumer(exchange ->
template.send(route, exchange));
+ }
+
+ private static MessageExt message() {
+ MessageExt messageExt = new MessageExt();
+ messageExt.setTopic("START_TOPIC");
+ messageExt.setBody("Hello".getBytes(StandardCharsets.UTF_8));
+ return messageExt;
+ }
+
+ @Override
+ protected RouteBuilder createRouteBuilder() {
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("direct:fail")
+ .throwException(new IllegalStateException("Forced"));
+
+ from("direct:rollback")
+ .markRollbackOnly();
+
+ from("direct:ok")
+ .convertBodyTo(String.class)
+ .to("mock:result");
+ }
+ };
+ }
+}
diff --git
a/components/camel-rocketmq/src/test/java/org/apache/camel/component/rocketmq/RocketMQProducerSendFailureTest.java
b/components/camel-rocketmq/src/test/java/org/apache/camel/component/rocketmq/RocketMQProducerSendFailureTest.java
new file mode 100644
index 000000000000..15504811ffdd
--- /dev/null
+++
b/components/camel-rocketmq/src/test/java/org/apache/camel/component/rocketmq/RocketMQProducerSendFailureTest.java
@@ -0,0 +1,48 @@
+/*
+ * 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.rocketmq;
+
+import org.apache.camel.Exchange;
+import org.apache.camel.ExchangePattern;
+import org.apache.camel.test.junit6.CamelTestSupport;
+import org.apache.rocketmq.client.exception.MQClientException;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+
+/**
+ * A send that fails asynchronously (here: the client rejects an empty message
before it contacts the name server) must
+ * fail the exchange, also for an InOut exchange without replyToTopic.
+ */
+public class RocketMQProducerSendFailureTest extends CamelTestSupport {
+
+ private static final String ROCKETMQ_URI =
"rocketmq:START_TOPIC?namesrvAddr=localhost:9876&producerGroup=p1";
+
+ @Test
+ public void testInOnlySendFailure() {
+ Exchange exchange = template.send(ROCKETMQ_URI,
ExchangePattern.InOnly, e -> e.getIn().setBody(""));
+
+ assertInstanceOf(MQClientException.class, exchange.getException());
+ }
+
+ @Test
+ public void testInOutSendFailureWithoutReplyToTopic() {
+ Exchange exchange = template.send(ROCKETMQ_URI, ExchangePattern.InOut,
e -> e.getIn().setBody(""));
+
+ assertInstanceOf(MQClientException.class, exchange.getException());
+ }
+}
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 f89769709b02..38f4fc054357 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
@@ -783,6 +783,15 @@ the message body as if the route had succeeded, which
usually echoed the request
body now replies with an empty (`null`) body instead of not replying, so the
sender no longer waits until its reply timeout (30 seconds by default).
+=== camel-rocketmq - a message whose route failed is consumed again
+
+The consumer acknowledged a message (`CONSUME_SUCCESS`) whenever the processor
returned, also when the route failed:
+a failed route sets the exception on the exchange instead of throwing it, so
the message was lost. The consumer now
+answers `RECONSUME_LATER` for a failed exchange (or one marked rollback only),
and RocketMQ redelivers the message with an increasing delay (up to
+the maximum reconsume times of the consumer group, 16 by default, then it goes
to the dead letter queue of the group).
+A route that fails for a message, and relied on the message being dropped, now
receives it again; handle the
+exception in the route (for example with `onException(...).handled(true)`) to
acknowledge the message anyway.
+
=== Components and Language removal
==== camel-csimple, camel-csimple-joor and csimple-maven-plugin