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

Reply via email to