davsclaus commented on code in PR #27173: URL: https://github.com/apache/camel/pull/27173#discussion_r4171970482
########## components/camel-sjms/src/main/java/org/apache/camel/component/sjms/consumer/BatchConsumerWorker.java: ########## @@ -0,0 +1,143 @@ +/* + * 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; + +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); + + BatchConsumerWorker(SjmsEndpoint endpoint, BatchEndpointMessageListener batchListener, + MessageConsumer consumer, Session session) { + this.endpoint = endpoint; + this.batchListener = batchListener; + this.consumer = consumer; + this.session = session; + } + + void shutdown() { + running.set(false); + } + + boolean isShutdownRequested() { + return !running.get(); + } + + private boolean isRedeliverable() { + return endpoint.isTransacted() + || endpoint.getAcknowledgementMode() == SessionAcknowledgementType.CLIENT_ACKNOWLEDGE; + } + + @Override + public void run() { + 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) { + dispatch(buffer); + buffer = new ArrayList<>(); + } + } + + if (!buffer.isEmpty()) { + dispatch(buffer); // graceful-stop drain + } Review Comment: When the container stops the worker because of a connection failure (`BatchMessageListenerContainer.onException` → `invalidateBatchWorkers()` → `worker.shutdown()`), the loop exits normally and this drain runs the buffered messages through the route on a dead session. The commit/ack then fails, so with transacted or `CLIENT_ACKNOWLEDGE` the broker redelivers them and they are processed twice. Please discard the buffer here when the stop comes from a connection failure, as the `JMSException` path below already does, and only drain on a graceful stop. ########## components/camel-sjms/src/main/java/org/apache/camel/component/sjms/consumer/BatchConsumerWorker.java: ########## @@ -0,0 +1,143 @@ +/* + * 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; + +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); + + BatchConsumerWorker(SjmsEndpoint endpoint, BatchEndpointMessageListener batchListener, + MessageConsumer consumer, Session session) { + this.endpoint = endpoint; + this.batchListener = batchListener; + this.consumer = consumer; + this.session = session; + } + + void shutdown() { + running.set(false); + } + + boolean isShutdownRequested() { + return !running.get(); + } + + private boolean isRedeliverable() { + return endpoint.isTransacted() + || endpoint.getAcknowledgementMode() == SessionAcknowledgementType.CLIENT_ACKNOWLEDGE; + } + + @Override + public void run() { + 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) { + dispatch(buffer); + buffer = new ArrayList<>(); + } + } + + if (!buffer.isEmpty()) { + dispatch(buffer); // graceful-stop drain + } + } catch (JMSException e) { + if (!buffer.isEmpty()) { + if (isRedeliverable()) { + LOG.error("Discarding {} buffered message(s) on {} after connection failure; " + + "unacknowledged/uncommitted, will be redelivered", + buffer.size(), endpoint.getEndpointUri()); + throw new BatchConsumerWorkerException(e); + } else { + dispatch(buffer); + LOG.warn("Connection failed on {} with {} already-acknowledged message(s) buffered; " + + "attempting best-effort dispatch since they cannot be redelivered", + endpoint.getEndpointUri(), buffer.size()); + } + } else { + throw new BatchConsumerWorkerException(e); + } + } + } + + private void dispatch(List<Message> buffer) { + try { + batchListener.onBatch(buffer, session); + } catch (Exception e) { + LOG.warn("Error dispatching batch of {} message(s) on {}", buffer.size(), + endpoint.getEndpointUri(), e); + } + } Review Comment: This only logs. Please hand the exception to `consumer.getExceptionHandler()` (as the non-batch consumer does), so `bridgeErrorHandler` and custom exception handlers work in batching mode. ########## components/camel-sjms/src/main/java/org/apache/camel/component/sjms/consumer/BatchEndpointMessageListener.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.consumer; + +import java.util.ArrayList; +import java.util.List; + +import jakarta.jms.Message; +import jakarta.jms.Session; + +import org.apache.camel.Exchange; +import org.apache.camel.Processor; +import org.apache.camel.component.sjms.SjmsConstants; +import org.apache.camel.component.sjms.SjmsConsumer; +import org.apache.camel.component.sjms.SjmsEndpoint; +import org.apache.camel.component.sjms.SjmsHelper; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +public class BatchEndpointMessageListener { + + private static final Logger LOG = LoggerFactory.getLogger(BatchEndpointMessageListener.class); + + private final SjmsConsumer consumer; + private final SjmsEndpoint endpoint; + private final Processor processor; + + public BatchEndpointMessageListener(SjmsConsumer consumer, SjmsEndpoint endpoint, Processor processor) { + this.consumer = consumer; + this.endpoint = endpoint; + this.processor = processor; + } + + private Exchange aggregate(List<Message> rawMessages, Session session) { + List<Exchange> exchanges = new ArrayList<Exchange>(rawMessages.size()); + for (Message m : rawMessages) { + Exchange e = endpoint.createExchange(m, session); + // Force eager materialization of JMS message headers and body into the + // Camel Exchange, before the session is committed/closed after dispatch. + e.getIn().getHeaders(); + e.getIn().getBody(); + + exchanges.add(e); + } + + Exchange batchExchange = consumer.createExchange(false); + batchExchange.getIn().setBody(exchanges); + batchExchange.setProperty(SjmsConstants.JMS_SESSION, session); + batchExchange.getMessage().setHeader(SjmsConstants.SJMS_BATCH_SIZE_HEADER, rawMessages.size()); + + return batchExchange; + } + + void onBatch(List<Message> rawMessages, Session session) throws Exception { + Exchange batchExchange = null; + Exception failure = null; + try { + batchExchange = aggregate(rawMessages, session); + processor.process(batchExchange); + } catch (Exception e) { + failure = e; + } + + boolean failed = failure != null || batchExchange.isFailed() || batchExchange.isRollbackOnly(); + Message lastMessage = rawMessages.get(rawMessages.size() - 1); + + if (!failed) { + SjmsHelper.commitIfNeeded(session, lastMessage); + } else { + Exception cause = failure != null ? failure : batchExchange.getException(); + if (cause != null) { + LOG.warn("Batch of {} message(s) failed processing on {}: {}", rawMessages.size(), + endpoint.getEndpointUri(), cause.getMessage(), cause); + } + SjmsHelper.rollbackIfNeeded(session); Review Comment: A rollback-only batch (`markRollbackOnly()`, no exception) is rolled back here without any log line. Please log it as well, and route a real failure to the consumer's `ExceptionHandler` rather than only logging a WARN. ########## components/camel-sjms/src/main/java/org/apache/camel/component/sjms/SjmsConstants.java: ########## @@ -38,5 +38,8 @@ public interface SjmsConstants { description = "Provides an explicit ReplyTo destination (overrides any incoming value of Message.getJMSReplyTo() in consumer)", javaType = "String") String JMS_REPLY_TO = JmsConstants.JMS_REPLY_TO; - + @Metadata(label = "consumer, batch", + description = "The size of the batch when using the batching consumer option.", + javaType = "String") Review Comment: The value is an int, and the space in the label makes the generated JSON group `" batch"`: ```suggestion @Metadata(label = "consumer,batch", description = "The size of the batch when using the batching consumer option.", javaType = "Integer") ``` ########## components/camel-sjms/src/main/java/org/apache/camel/component/sjms/SjmsEndpoint.java: ########## @@ -449,6 +488,35 @@ 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) { + if (getExchangePattern().isOutCapable()) { + throw new IllegalArgumentException("batchInterval must be 0 or greater."); + } + } Review Comment: This check can never fire: it is nested inside `isOutCapable()`, which already threw at L499. Also, `batchSize <= 0` is not rejected; with `batchInterval=0` a batch then never completes until the consumer stops. ```suggestion if (getBatchInterval() < 0) { throw new IllegalArgumentException("batchInterval must be 0 or greater."); } if (getBatchSize() <= 0) { throw new IllegalArgumentException("batchSize must be greater than 0."); } ``` -- 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]
