gnodet commented on code in PR #27173:
URL: https://github.com/apache/camel/pull/27173#discussion_r4239054545


##########
catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/sjms-component.adoc:
##########
@@ -312,3 +312,34 @@ Currently, the only correlation strategy is to use the 
`JMSCorrelationId`.
 The _InOut_ Consumer uses this strategy as well ensuring that all
 response messages to the included `JMSReplyTo` destination also have the
 `JMSCorrelationId` copied from the request as well.
+
+=== Batch consuming
+
+The consumer can be configured to receive messages in batches instead of one 
at a time, by setting
+`batching=true`. Instead of one Exchange per message, the route receives a 
single Exchange whose body
+is a `List<Exchange>`, one per JMS message in the batch, with a 
`CamelSjmsBatchSize` header giving the
+batch's size. Each Exchange in the batch is comparable to the Exchange 
generated by a non-batching consumer.
+
+Batching only supports the InOnly exchange pattern. There is no 
reply-to/request-response support for
+batched consumption — a batch of N unrelated messages has no well-defined 
single reply, so InOut
+routes are not supported when `batching=true`.
+
+A batch completes when either `batchSize` messages have been received, or the 
`batchInterval` —
+measured from the first message in the batch - has elapsed.
+
+Commit/acknowledgement happens once per batch rather than once per message. 
With
+`acknowledgementMode=AUTO_ACKNOWLEDGE`, messages are acknowledged as they're 
received into the batch,
+before routing — so a failed batch cannot be redelivered. Use 
`transacted=true` or
+`acknowledgementMode=CLIENT_ACKNOWLEDGE` if the whole batch should be 
redelivered together on failure.
+
+A failed batch is rolled back and redelivered as a whole, so the route should 
be idempotent. One bad message
+fails, and redelivers, the whole batch: all its messages share the delivery 
count, and they end up together
+in the broker's dead letter queue (if one is configured). How often and how 
fast the broker redelivers is
+up to its redelivery policy, for example `max-delivery-attempts` and 
`redelivery-delay` in ActiveMQ Artemis.
+To avoid this, handle failures per message inside the batch, for example with 
`split(body())` and `doTry`,
+or with `onException(...).handled(true)`. The messages of a redelivered batch 
have the `JMSRedelivered` header set.
+
+With `concurrentConsumers` greater than 1, each consumer collects and routes 
its own batches, so several
+batches are processed in parallel.
+
+The batch consumer never sends a reply, so a `JMSReplyTo` on the incoming 
messages is ignored.

Review Comment:
   💬 **Missing final newline.** The catalog copy of `sjms-component.adoc` is 
missing its trailing newline (the source file at 
`components/camel-sjms/src/main/docs/sjms-component.adoc` has one). Add a 
newline after the last line of the file.



##########
components/camel-sjms2/src/test/java/org/apache/camel/component/sjms2/batch/BatchConsumerTransactedTest.java:
##########
@@ -0,0 +1,175 @@
+/*
+ * 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.sjms2.batch;
+
+import java.util.List;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import jakarta.jms.ConnectionFactory;
+
+import org.apache.camel.Exchange;
+import org.apache.camel.Processor;
+import org.apache.camel.RoutesBuilder;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.apache.camel.component.sjms2.support.Jms2TestSupport;
+import org.apache.camel.test.infra.artemis.services.ArtemisService;
+import org.apache.camel.test.infra.artemis.services.ArtemisServiceFactory;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.RegisterExtension;
+
+import static java.lang.String.format;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+
+public class BatchConsumerTransactedTest extends Jms2TestSupport {
+
+    private static final String QUEUE_NAME_TEMPLATE = 
"sjms2.batch.consumer.%s.BatchConsumerTransactedTest";
+
+    static final String MOCK_START = "mock:%s.start";
+    static final String MOCK_FINISH = "mock:%s.complete";
+
+    private static final String ROUTE_ID_SESSION_TX = "tx";
+    private static final String ROUTE_ID_CLIENT_ACK_NO_TX = "no-tx-client-ack";
+    private static final String ROUTE_ID_AUTO_ACK_NO_TX = "no-tx-auto";
+
+    private static final String MESSAGE_TEXT = "Message %d";
+
+    @RegisterExtension
+    public static ArtemisService service = 
ArtemisServiceFactory.createTCPAllProtocolsService();
+
+    @Test
+    public void testClientAcknowledgedNotTransacted() throws Exception {
+        MockEndpoint mockStart = getMockEndpoint(format(MOCK_START, 
ROUTE_ID_CLIENT_ACK_NO_TX));
+        mockStart.expectedMessageCount(2);
+
+        MockEndpoint mockFinish = getMockEndpoint(format(MOCK_FINISH, 
ROUTE_ID_CLIENT_ACK_NO_TX));
+        mockFinish.expectedMessageCount(1);
+
+        for (int i = 1; i <= 5; i++) {
+            template.sendBody("sjms2:queue:" + format(QUEUE_NAME_TEMPLATE, 
ROUTE_ID_CLIENT_ACK_NO_TX),
+                    format(MESSAGE_TEXT, i));
+        }
+
+        MockEndpoint.assertIsSatisfied(context);
+
+        List<Exchange> batch = 
mockFinish.getExchanges().get(0).getMessage().getBody(List.class);
+        assertNotNull(batch);
+        assertEquals(5, batch.size());
+
+        int messageNumber = 1;
+        for (Exchange batchExchange : batch) {
+            assertEquals(
+                    String.format(MESSAGE_TEXT, messageNumber++),
+                    batchExchange.getMessage().getBody(String.class));
+        }
+    }
+
+    @Test
+    public void testSessionTransacted() throws Exception {
+        MockEndpoint mockStart = getMockEndpoint(format(MOCK_START, 
ROUTE_ID_SESSION_TX));
+        mockStart.expectedMessageCount(2);
+
+        MockEndpoint mockFinish = getMockEndpoint(format(MOCK_FINISH, 
ROUTE_ID_SESSION_TX));
+        mockFinish.expectedMessageCount(1);
+
+        for (int i = 1; i <= 5; i++) {
+            template.sendBody("sjms2:queue:" + format(QUEUE_NAME_TEMPLATE, 
ROUTE_ID_SESSION_TX),
+                    format(MESSAGE_TEXT, i));
+        }
+
+        MockEndpoint.assertIsSatisfied(context);
+
+        List<Exchange> batch = 
mockFinish.getExchanges().get(0).getMessage().getBody(List.class);
+        assertNotNull(batch);
+        assertEquals(5, batch.size());
+
+        int messageNumber = 1;
+        for (Exchange batchExchange : batch) {
+            assertEquals(
+                    String.format(MESSAGE_TEXT, messageNumber++),
+                    batchExchange.getMessage().getBody(String.class));
+        }
+    }
+
+    @Test
+    public void testAutoAcknowledgedNotTransacted() throws Exception {
+        MockEndpoint mockStart = getMockEndpoint(format(MOCK_START, 
ROUTE_ID_AUTO_ACK_NO_TX));
+        mockStart.expectedMessageCount(1);
+
+        MockEndpoint mockFinish = getMockEndpoint(format(MOCK_FINISH, 
ROUTE_ID_AUTO_ACK_NO_TX));
+        mockFinish.expectedMessageCount(0);
+
+        for (int i = 1; i <= 5; i++) {
+            template.sendBody("sjms2:queue:" + format(QUEUE_NAME_TEMPLATE, 
ROUTE_ID_AUTO_ACK_NO_TX),
+                    format(MESSAGE_TEXT, i));
+        }
+
+        MockEndpoint.assertIsSatisfied(context);
+    }
+
+    @Override
+    protected RoutesBuilder createRouteBuilder() {
+        return new RouteBuilder() {
+            @Override
+            public void configure() {
+                
fromF("sjms2:%s?batching=true&batchSize=5&batchInterval=10000&transacted=%s&acknowledgementMode=%s",
+                        format(QUEUE_NAME_TEMPLATE, ROUTE_ID_AUTO_ACK_NO_TX),
+                        false,
+                        "AUTO_ACKNOWLEDGE")
+                        .routeId(ROUTE_ID_AUTO_ACK_NO_TX)
+                        .toF(MOCK_START, ROUTE_ID_AUTO_ACK_NO_TX)
+                        .process(new ThrowExceptionProcessor())
+                        .toF(MOCK_FINISH, ROUTE_ID_AUTO_ACK_NO_TX);
+
+                
fromF("sjms2:%s?batching=true&batchSize=5&batchInterval=10000&transacted=%s&acknowledgementMode=%s",
+                        format(QUEUE_NAME_TEMPLATE, ROUTE_ID_SESSION_TX),
+                        true,
+                        "SESSION_TRANSACTED")
+                        .routeId(ROUTE_ID_SESSION_TX)
+                        .toF(MOCK_START, ROUTE_ID_SESSION_TX)
+                        .process(new ThrowExceptionProcessor())
+                        .toF(MOCK_FINISH, ROUTE_ID_SESSION_TX);
+
+                
fromF("sjms2:%s?batching=true&batchSize=5&batchInterval=10000&transacted=%s&acknowledgementMode=%s",
+                        format(QUEUE_NAME_TEMPLATE, ROUTE_ID_CLIENT_ACK_NO_TX),
+                        false,
+                        "CLIENT_ACKNOWLEDGE")
+                        .routeId(ROUTE_ID_CLIENT_ACK_NO_TX)
+                        .toF(MOCK_START, ROUTE_ID_CLIENT_ACK_NO_TX)
+                        .process(new ThrowExceptionProcessor())
+                        .toF(MOCK_FINISH, ROUTE_ID_CLIENT_ACK_NO_TX);
+            }
+        };
+    }
+
+    protected ConnectionFactory getConnectionFactory() throws Exception {
+        return getConnectionFactory(service.serviceAddress());
+    }
+
+    private static class ThrowExceptionProcessor implements Processor {
+        private final AtomicInteger counter = new AtomicInteger();
+
+        @Override
+        public void process(org.apache.camel.Exchange exchange) {

Review Comment:
   💬 **Redundant FQN — `Exchange` is already imported.** Use the short name:
   
   ```suggestion
           public void process(Exchange exchange) {
   ```



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to