davsclaus commented on code in PR #27173:
URL: https://github.com/apache/camel/pull/27173#discussion_r4220145710
##########
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:
Small wording fixes still open: the em dash, 'postive' → 'positive', and
maybe 'Set to 0 to disable the interval' instead of 'unlimited'. The generated
files need regenerating afterwards.
##########
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:
Thanks for checking. You're right that `SimpleMessageListenerContainer`
(line 162) does the same for the non-batch consumer, so this matches existing
behaviour. Fine to leave it as is.
##########
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();
Review Comment:
`stopping` is set in `doStop` (line 94) but only cleared in the constructor.
After a route stop/start, `onWorkerExit` sees `stopping == true`, so a worker
that dies on a receive `JMSException` no longer triggers recovery. Could you
reset it here, before `super.doStart()`?
##########
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:
```suggestion
throw new IllegalArgumentException("batchSize must be greater
than 0");
```
##########
components/camel-sjms/src/test/java/org/apache/camel/component/sjms/batch/BatchConsumerRollbackOnlyTest.java:
##########
@@ -0,0 +1,92 @@
+/*
+ * 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.RoutesBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.apache.camel.component.sjms.support.JmsTestSupport;
+import org.junit.jupiter.api.Test;
+
+import static java.lang.String.format;
+import static
org.apache.camel.component.sjms.batch.BatchTestHelper.BATCH_ROUTEBUILDER_MOCK_FINISH;
+import static
org.apache.camel.component.sjms.batch.BatchTestHelper.BATCH_ROUTEBUILDER_MOCK_START;
+import static
org.apache.camel.component.sjms.batch.BatchTestHelper.assertBatchSizesInOrder;
+import static
org.apache.camel.component.sjms.batch.BatchTestHelper.createRoute;
+import static
org.apache.camel.component.sjms.batch.BatchTestHelper.sendMessages;
+
+public class BatchConsumerRollbackOnlyTest extends JmsTestSupport {
+
+ private static final String QUEUE_NAME_TEMPLATE =
"batch.consumer.%s.BatchConsumerRollbackOnlyTest";
+
+ private static final String ROUTE_ID_SESSION_TX = "tx";
+ private static final String ROUTE_ID_CLIENT_ACK_NO_TX = "no-tx-client-ack";
+
+ @Test
+ public void testSessionTransacted() throws Exception {
+ MockEndpoint mockStart =
getMockEndpoint(format(BATCH_ROUTEBUILDER_MOCK_START, ROUTE_ID_SESSION_TX));
+ mockStart.expectedMessageCount(2);
+
+ MockEndpoint mockFinish =
getMockEndpoint(format(BATCH_ROUTEBUILDER_MOCK_FINISH, ROUTE_ID_SESSION_TX));
+ mockFinish.expectedMessageCount(1);
+
+ sendMessages(template, format("sjms:queue:" + QUEUE_NAME_TEMPLATE,
ROUTE_ID_SESSION_TX),
+ 5);
+
+ MockEndpoint.assertIsSatisfied(context);
+ assertBatchSizesInOrder(mockFinish, 5);
+ }
+
+ @Test
+ public void testClientAcknowledgedNotTransacted() throws Exception {
+ MockEndpoint mockStart =
getMockEndpoint(format(BATCH_ROUTEBUILDER_MOCK_START,
ROUTE_ID_CLIENT_ACK_NO_TX));
+ mockStart.expectedMessageCount(2);
+
+ MockEndpoint mockFinish =
getMockEndpoint(format(BATCH_ROUTEBUILDER_MOCK_FINISH,
ROUTE_ID_CLIENT_ACK_NO_TX));
+ mockFinish.expectedMessageCount(1);
+
+ sendMessages(template,
+ format("sjms:queue:" + QUEUE_NAME_TEMPLATE,
ROUTE_ID_CLIENT_ACK_NO_TX),
+ 5);
+
+ MockEndpoint.assertIsSatisfied(context);
+ assertBatchSizesInOrder(mockFinish, 5);
+ }
+
+ @Override
+ protected RoutesBuilder[] createRouteBuilders() {
+ return new org.apache.camel.RoutesBuilder[] {
+ createRoute(QUEUE_NAME_TEMPLATE, ROUTE_ID_SESSION_TX, true, 5,
+ 1000,
+ true, null, 1, new MarkRollBackProcessor()),
+ createRoute(QUEUE_NAME_TEMPLATE, ROUTE_ID_CLIENT_ACK_NO_TX,
+ true,
+ 5, 1000, false, "CLIENT_ACKNOWLEDGE", 1, new
MarkRollBackProcessor())
+ };
+ }
+
+ private static class MarkRollBackProcessor implements
org.apache.camel.Processor {
+ private final java.util.concurrent.atomic.AtomicInteger counter = new
java.util.concurrent.atomic.AtomicInteger();
+
+ @Override
+ public void process(org.apache.camel.Exchange exchange) {
Review Comment:
Please import `Processor`, `Exchange` and `AtomicInteger` (and use the
already-imported `RoutesBuilder` at line 71) instead of fully qualified names,
per CLAUDE.md. The OpenRewrite step in the build would rewrite these anyway.
--
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]