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


##########
components/camel-sjms/src/main/java/org/apache/camel/component/sjms/consumer/BatchConsumerWorker.java:
##########
@@ -0,0 +1,178 @@
+/*
+ * 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.sjms.consumer;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+
+import jakarta.jms.JMSException;
+import jakarta.jms.Message;
+import jakarta.jms.MessageConsumer;
+import jakarta.jms.Session;
+
+import org.apache.camel.component.sjms.SjmsEndpoint;
+import org.apache.camel.component.sjms.jms.SessionAcknowledgementType;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import static org.apache.camel.component.sjms.SjmsHelper.commitIfNeeded;
+import static org.apache.camel.component.sjms.SjmsHelper.rollbackIfNeeded;
+
+public class BatchConsumerWorker implements Runnable {
+
+    private static final Logger LOG = 
LoggerFactory.getLogger(BatchConsumerWorker.class);
+
+    private static final long JMS_CONSUMER_RECEIVE_WAKE_INTERVAL_TIMEOUT = 
1000L;
+    private static final long JMS_CONSUMER_RECEIVE_MIN_TIMEOUT = 100L;
+
+    private final SjmsEndpoint endpoint;
+    private final BatchEndpointMessageListener batchListener;
+    private final MessageConsumer consumer;
+    private final Session session;
+    private final AtomicBoolean running = new AtomicBoolean(true);
+    private volatile boolean jmsUnhealthy;
+
+    BatchConsumerWorker(SjmsEndpoint endpoint, BatchEndpointMessageListener 
batchListener,
+                        MessageConsumer consumer, Session session) {
+        this.endpoint = endpoint;
+        this.batchListener = batchListener;
+        this.consumer = consumer;
+        this.session = session;
+    }
+
+    void shutdown() {
+        LOG.debug("Shutdown requested");
+        running.set(false);
+    }
+
+    boolean isShutdownRequested() {
+        return !running.get();
+    }
+
+    void invalidate() {
+        jmsUnhealthy = true;
+        shutdown();
+    }
+
+    private boolean isRedeliverable() {
+        return endpoint.isTransacted()
+                || endpoint.getAcknowledgementMode() == 
SessionAcknowledgementType.CLIENT_ACKNOWLEDGE;
+    }
+
+    @Override
+    public void run() {
+        LOG.debug("run START");
+        int batchSize = endpoint.getBatchSize();
+        long batchInterval = endpoint.getBatchInterval();
+        List<Message> buffer = new ArrayList<>();
+        long batchStartTime = 0L;
+
+        try {
+            while (running.get()) {
+                long waitMillis
+                        = computeWaitMillis(buffer.isEmpty(), batchStartTime, 
batchInterval);
+
+                Message msg = consumer.receive(waitMillis);
+
+                if (msg != null) {
+                    if (buffer.isEmpty()) {
+                        batchStartTime = System.nanoTime();
+                    }
+                    buffer.add(msg);
+                }
+
+                boolean sizeReached = batchSize > 0 && buffer.size() >= 
batchSize;
+                boolean intervalElapsed = batchInterval > 0 && 
!buffer.isEmpty()
+                        && TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - 
batchStartTime) >= batchInterval;
+
+                if (sizeReached || intervalElapsed) {
+                    dispatchOrDiscard(buffer);
+                    buffer = new ArrayList<>();
+                }
+            }
+
+            if (!buffer.isEmpty()) {
+                dispatchOrDiscard(buffer); // graceful-stop drain
+            }
+        } catch (JMSException e) {
+            jmsUnhealthy = true;
+            if (!buffer.isEmpty()) {
+                dispatchOrDiscard(buffer);
+            }
+            throw new BatchConsumerWorkerException(e);
+        }
+        LOG.info("run STOP");
+    }
+
+    private void dispatchOrDiscard(List<Message> batch) {
+        if (jmsUnhealthy) {
+            if (isRedeliverable()) {
+                LOG.warn("JMS is unhealthy on {}: discarding {} buffered 
message(s); "
+                         + "unacknowledged/uncommitted, will be redelivered",
+                        endpoint.getEndpointUri(), batch.size());
+                return;
+            }
+            LOG.warn("JMS is unhealthy on {}: attempting best-effort dispatch 
of {} already-acknowledged "
+                     + "message(s) since they cannot be redelivered",
+                    endpoint.getEndpointUri(), batch.size());
+        }
+        onBatch(batch);
+    }
+
+    private void onBatch(List<Message> batch) {
+        try {
+            doOnBatch(batch);
+        } catch (Exception e) {
+            if (e instanceof JMSException jmsException) {
+                if (endpoint.getExceptionListener() != null) {
+                    endpoint.getExceptionListener().onException(jmsException);
+                }
+            } else {
+                LOG.warn("Execution of JMS message listener failed. This 
exception is ignored.", e);

Review Comment:
   ⚠️ **Still present (PRRT_kwDOAAMl7c6opow1, PRRT_kwDOAAMl7c6p9pZG): Non-JMS 
exceptions are silently swallowed.** When `doOnBatch()` throws a non-JMS 
exception (e.g. a route failure, processor exception, or 
`RuntimeCamelException` from `onBatch()`), the `else` branch logs a WARN and 
discards it. The Camel convention — used by the non-batching 
`EndpointMessageListener` — is to delegate to 
`consumer.getExceptionHandler().handleException(...)`, which makes 
dead-letter-channel, error-handler, and monitoring integrations work as users 
expect.
   
   `BatchConsumerWorker` already holds a reference to `batchListener`, which 
holds `consumer`. A small delegating method on `BatchEndpointMessageListener`:
   ```java
   void handleException(String message, Throwable t) {
       consumer.getExceptionHandler().handleException(message, t);
   }
   ```
   Then in the `else` branch:
   ```suggestion
               } else {
                   batchListener.handleException(
                           "Execution of JMS batch message listener failed", e);
               }
   ```



##########
components/camel-sjms/src/main/java/org/apache/camel/component/sjms/SjmsEndpoint.java:
##########
@@ -285,6 +291,18 @@ public class SjmsEndpoint extends DefaultEndpoint
     @UriParam(defaultValue = "false", label = "advanced",
               description = "Sets whether synchronous processing should be 
strictly used")
     private boolean synchronous;
+    @UriParam(label = "consumer,batch", defaultValue = "false",
+              description = "Enable batch consuming. The route receives one 
Exchange per batch, whose body"
+                            + " is a List<Exchange> of the individual JMS 
messages, instead of one Exchange per message.")
+    private boolean batching;
+    @UriParam(defaultValue = "100", label = "consumer,batch",
+              description = "Maximum number of messages per batch.")
+    private int batchSize = 100;
+    @UriParam(defaultValue = "1000", label = "consumer,batch", javaType = 
"java.time.Duration",
+              description = "Time in millis, measured from the first message 
received into a new batch, after which "
+                            + "the batch is dispatched even if batchSize has 
not been reached — comparable to the Aggregator "
+                            + "EIP's completionInterval. Default is 1000 ms, 
that is 1 second. Interval should be a postive value. Set to 0 for unlimited 
(not recommended).")

Review Comment:
   Nit: "postive" → "positive". This typo propagates into the generated catalog 
JSON and DSL Javadoc for both sjms and sjms2.
   ```suggestion
                               + "EIP's completionInterval. Default is 1000 ms, 
that is 1 second. Interval should be a positive value. Set to 0 for unlimited 
(not recommended.)")
   ```



##########
components/camel-sjms/src/main/java/org/apache/camel/component/sjms/SjmsEndpoint.java:
##########
@@ -449,6 +488,37 @@ public MessageListenerContainer 
createMessageListenerContainer(SjmsEndpoint endp
         return answer;
     }
 
+    public BatchMessageListenerContainer createBatchMessageListenerContainer(
+            SjmsEndpoint endpoint) {
+        BatchMessageListenerContainer answer = new 
BatchMessageListenerContainer(endpoint);
+        answer.setConcurrentConsumers(concurrentConsumers);
+        return answer;
+    }
+
+    private void validateBatchingOptions() {
+        if (getExchangePattern().isOutCapable()) {
+            throw new IllegalArgumentException("SjmsConsumer does not support 
exchangePattern=InOut in batching mode");
+        }
+
+        if (getBatchInterval() < 0) {
+            throw new IllegalArgumentException("batchInterval must be 0 or 
greater.");
+        }
+
+        if (getBatchSize() <= 0) {
+            throw new IllegalArgumentException("batchSize must greater than 
0");

Review Comment:
   Nit: missing "be" — "must greater" → "must be greater".
   ```suggestion
               throw new IllegalArgumentException("batchSize must be greater 
than 0");
   ```



##########
components/camel-sjms/src/test/java/org/apache/camel/component/sjms/batch/BatchConsumerValidationTest.java:
##########
@@ -0,0 +1,89 @@
+/*
+ * 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.sjms.batch;
+
+import org.apache.camel.Consumer;
+import org.apache.camel.component.sjms.SjmsEndpoint;
+import org.apache.camel.component.sjms.support.JmsTestSupport;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.ValueSource;
+
+import static org.junit.jupiter.api.Assertions.*;

Review Comment:
   Nit: `import static org.junit.jupiter.api.Assertions.*` — use explicit 
static imports per Camel convention (e.g. `assertEquals`, `assertThrows`, 
`assertTrue`, `assertNotNull`).



##########
components/camel-sjms/src/main/java/org/apache/camel/component/sjms/consumer/BatchMessageListenerContainer.java:
##########
@@ -0,0 +1,156 @@
+/*
+ * 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.sjms.consumer;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.locks.ReentrantLock;
+
+import jakarta.jms.JMSException;
+import jakarta.jms.MessageConsumer;
+import jakarta.jms.Session;
+
+import org.apache.camel.component.sjms.SjmsEndpoint;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+public class BatchMessageListenerContainer extends 
SimpleMessageListenerContainer {
+
+    private static final Logger LOG = 
LoggerFactory.getLogger(BatchMessageListenerContainer.class);
+
+    private final SjmsEndpoint endpoint;
+    private BatchEndpointMessageListener batchListener;
+    private ExecutorService workerExecutorService;
+    private final ReentrantLock workersLock = new ReentrantLock();
+    private final List<BatchConsumerWorker> workers = new ArrayList<>();
+    private volatile boolean stopping;
+
+    public BatchMessageListenerContainer(SjmsEndpoint endpoint) {
+        super(endpoint);
+        this.endpoint = endpoint;
+        this.stopping = false;
+    }
+
+    public void setBatchListener(BatchEndpointMessageListener batchListener) {
+        this.batchListener = batchListener;
+    }
+
+    @Override
+    protected void doStart() throws Exception {
+        workerExecutorService = 
endpoint.getCamelContext().getExecutorServiceManager().newFixedThreadPool(
+                this, "SjmsBatchConsumer[" + endpoint.getDestinationName() + 
"]",
+                Math.max(1, this.getConcurrentConsumers()));
+
+        // triggers connection + session/consumer creation, calling 
configureConsumer() below
+        // for each session per concurrentConsumers, and re-invokes it again 
on reconnection
+        super.doStart();
+    }
+
+    private void addWorker(BatchConsumerWorker worker) {
+        workersLock.lock();
+        try {
+            workers.add(worker);
+        } finally {
+            workersLock.unlock();
+        }
+    }
+
+    private boolean removeWorker(BatchConsumerWorker worker) {
+        workersLock.lock();
+        try {
+            return workers.remove(worker);
+        } finally {
+            workersLock.unlock();
+        }
+    }
+
+    @Override
+    protected void configureConsumer(MessageConsumer consumer, Session 
session) {
+        BatchConsumerWorker worker = new BatchConsumerWorker(
+                endpoint, batchListener, consumer, session);
+        addWorker(worker);
+        CompletableFuture.runAsync(worker, workerExecutorService)
+                .whenComplete((v, ex) -> onWorkerExit(worker, ex));
+    }
+
+    @Override
+    protected void doStop() throws Exception {
+        this.stopping = true;
+        shutdownWorkers();
+        // shutdown scheduled executor after all in-flight exchanges have 
completed
+        if (workerExecutorService != null) {
+            
getEndpoint().getCamelContext().getExecutorServiceManager().shutdownGraceful(workerExecutorService);
+            LOG.debug("BatchConsumerWorker excecutor service shutdown");

Review Comment:
   Nit: "excecutor" → "executor".
   ```suggestion
               LOG.debug("BatchConsumerWorker executor service shutdown");
   ```



##########
components/camel-sjms/src/test/java/org/apache/camel/component/sjms/batch/BatchTestHelper.java:
##########
@@ -0,0 +1,243 @@
+/*
+ * 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.sjms.batch;
+
+import java.lang.reflect.Field;
+import java.util.*;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.stream.*;
+
+import jakarta.jms.Connection;
+import jakarta.jms.ConnectionFactory;
+import jakarta.jms.ExceptionListener;
+import jakarta.jms.JMSContext;
+import jakarta.jms.JMSException;
+
+import org.apache.activemq.artemis.core.server.ActiveMQServer;
+import org.apache.activemq.artemis.core.server.ServerConsumer;
+import org.apache.activemq.artemis.core.server.ServerSession;
+import org.apache.camel.*;
+import org.apache.camel.builder.*;

Review Comment:
   Nit: Star imports (`java.util.*`, `java.util.stream.*`, 
`org.apache.camel.*`, `org.apache.camel.builder.*`) violate Camel's checkstyle 
rules. Replace with explicit imports for the specific classes used.



-- 
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