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 dd1d7cafed24 CAMEL-24196: camel-aws2-sqs - Fix QueueDoesNotExist 
swallowed, thread pool leak, and batch chunking
dd1d7cafed24 is described below

commit dd1d7cafed246bc8bcffaa0eeff133452234b964
Author: Claus Ibsen <[email protected]>
AuthorDate: Fri Jul 17 22:41:29 2026 +0200

    CAMEL-24196: camel-aws2-sqs - Fix QueueDoesNotExist swallowed, thread pool 
leak, and batch chunking
    
    Co-Authored-By: Claude Opus 4.6 <[email protected]>
---
 .../camel/component/aws2/sqs/Sqs2Consumer.java     | 26 +++++++++++
 .../camel/component/aws2/sqs/Sqs2Producer.java     | 53 +++++++++++++---------
 2 files changed, 57 insertions(+), 22 deletions(-)

diff --git 
a/components/camel-aws/camel-aws2-sqs/src/main/java/org/apache/camel/component/aws2/sqs/Sqs2Consumer.java
 
b/components/camel-aws/camel-aws2-sqs/src/main/java/org/apache/camel/component/aws2/sqs/Sqs2Consumer.java
index 1744dbfc96b5..e51d2912a363 100644
--- 
a/components/camel-aws/camel-aws2-sqs/src/main/java/org/apache/camel/component/aws2/sqs/Sqs2Consumer.java
+++ 
b/components/camel-aws/camel-aws2-sqs/src/main/java/org/apache/camel/component/aws2/sqs/Sqs2Consumer.java
@@ -338,6 +338,28 @@ public class Sqs2Consumer extends 
ScheduledBatchPollingConsumer {
         super.doStart();
     }
 
+    @Override
+    protected void doStop() throws Exception {
+        super.doStop();
+
+        if (ObjectHelper.isNotEmpty(timeoutExtender)) {
+            timeoutExtender.cancel();
+            timeoutExtender = null;
+        }
+        if (ObjectHelper.isNotEmpty(scheduledFuture)) {
+            scheduledFuture.cancel(true);
+            scheduledFuture = null;
+        }
+        if (ObjectHelper.isNotEmpty(scheduledExecutor)) {
+            
getEndpoint().getCamelContext().getExecutorServiceManager().shutdown(scheduledExecutor);
+            scheduledExecutor = null;
+        }
+        if (ObjectHelper.isNotEmpty(pollingTask)) {
+            pollingTask.close();
+            pollingTask = null;
+        }
+    }
+
     @Override
     protected void doShutdown() throws Exception {
         if (ObjectHelper.isNotEmpty(timeoutExtender)) {
@@ -594,6 +616,10 @@ public class Sqs2Consumer extends 
ScheduledBatchPollingConsumer {
 
             final PollingContext context = new PollingContext();
             final List<software.amazon.awssdk.services.sqs.model.Message> 
messages = poll(context);
+            if (context.isQueueMissing() && context.hasErrors()) {
+                context.rethrowIfFirstErrorIsRuntimeException();
+                throw new IOException("Queue %s does not 
exist".formatted(queueName), context.firstError());
+            }
             if (context.errorCount() == numberOfRequestsPerPoll) {
                 if (context.errorCount() == 1) {
                     context.rethrowIfFirstErrorIsRuntimeException();
diff --git 
a/components/camel-aws/camel-aws2-sqs/src/main/java/org/apache/camel/component/aws2/sqs/Sqs2Producer.java
 
b/components/camel-aws/camel-aws2-sqs/src/main/java/org/apache/camel/component/aws2/sqs/Sqs2Producer.java
index 1295bf948d05..a911d04871bb 100644
--- 
a/components/camel-aws/camel-aws2-sqs/src/main/java/org/apache/camel/component/aws2/sqs/Sqs2Producer.java
+++ 
b/components/camel-aws/camel-aws2-sqs/src/main/java/org/apache/camel/component/aws2/sqs/Sqs2Producer.java
@@ -17,8 +17,8 @@
 package org.apache.camel.component.aws2.sqs;
 
 import java.util.ArrayList;
-import java.util.Collection;
 import java.util.HashMap;
+import java.util.List;
 import java.util.Map;
 import java.util.Map.Entry;
 import java.util.UUID;
@@ -56,6 +56,7 @@ public class Sqs2Producer extends DefaultProducer {
     private static final Logger LOG = 
LoggerFactory.getLogger(Sqs2Producer.class);
 
     private static final int MAX_ATTRIBUTES = 10;
+    private static final int MAX_BATCH_ENTRIES = 10;
     private static final String MAX_MESSAGE
             = "Number of message headers exceeded. At most " + MAX_ATTRIBUTES 
+ " headers is allowed when sending to AWS SQS.";
 
@@ -121,32 +122,24 @@ public class Sqs2Producer extends DefaultProducer {
                 && 
ObjectHelper.isEmpty(getEndpoint().getConfiguration().getMessageGroupIdStrategy()))
 {
             throw new IllegalArgumentException("messageGroupIdStrategy must be 
set for FIFO queues.");
         }
-        SendMessageBatchRequest.Builder request = 
SendMessageBatchRequest.builder().queueUrl(getQueueUrl());
-        Collection<SendMessageBatchRequestEntry> entries = new ArrayList<>();
         if (exchange.getIn().getBody() instanceof Iterable) {
-            Iterable c = exchange.getIn().getBody(Iterable.class);
+            Iterable<?> c = exchange.getIn().getBody(Iterable.class);
+            List<SendMessageBatchRequestEntry> entries = new ArrayList<>();
             int index = 0;
             for (Object o : c) {
-                String object = (String) o;
                 SendMessageBatchRequestEntry.Builder entry = 
SendMessageBatchRequestEntry.builder();
                 entry.id(UUID.randomUUID().toString());
                 
entry.messageAttributes(translateAttributes(exchange.getIn().getHeaders(), 
exchange));
-                entry.messageBody(object);
+                entry.messageBody((String) o);
                 addDelay(entry, exchange);
                 configureFifoAttributes(entry, exchange, index++);
                 entries.add(entry.build());
             }
-            request.entries(entries);
-            SendMessageBatchResponse result = 
amazonSQS.sendMessageBatch(request.build());
-            Message message = getMessageForResponse(exchange);
-            message.setBody(result);
-            message.setHeader(Sqs2Constants.FAILED_MESSAGE_COUNT,
-                    ObjectHelper.isNotEmpty(result.failed()) ? 
result.failed().size() : 0);
-            message.setHeader(Sqs2Constants.SUCCESSFUL_MESSAGE_COUNT,
-                    ObjectHelper.isNotEmpty(result.successful()) ? 
result.successful().size() : 0);
+            sendBatchEntries(amazonSQS, exchange, entries);
         } else if (exchange.getIn().getBody() instanceof String) {
             String c = exchange.getIn().getBody(String.class);
             String[] elements = 
c.split(getConfiguration().getBatchSeparator());
+            List<SendMessageBatchRequestEntry> entries = new ArrayList<>();
             int index = 0;
             for (String o : elements) {
                 SendMessageBatchRequestEntry.Builder entry = 
SendMessageBatchRequestEntry.builder();
@@ -157,14 +150,7 @@ public class Sqs2Producer extends DefaultProducer {
                 configureFifoAttributes(entry, exchange, index++);
                 entries.add(entry.build());
             }
-            request.entries(entries);
-            SendMessageBatchResponse result = 
amazonSQS.sendMessageBatch(request.build());
-            Message message = getMessageForResponse(exchange);
-            message.setBody(result);
-            message.setHeader(Sqs2Constants.FAILED_MESSAGE_COUNT,
-                    ObjectHelper.isNotEmpty(result.failed()) ? 
result.failed().size() : 0);
-            message.setHeader(Sqs2Constants.SUCCESSFUL_MESSAGE_COUNT,
-                    ObjectHelper.isNotEmpty(result.successful()) ? 
result.successful().size() : 0);
+            sendBatchEntries(amazonSQS, exchange, entries);
         } else {
             SendMessageBatchRequest req = 
exchange.getIn().getBody(SendMessageBatchRequest.class);
             SendMessageBatchResponse result = amazonSQS.sendMessageBatch(req);
@@ -177,6 +163,29 @@ public class Sqs2Producer extends DefaultProducer {
         }
     }
 
+    private void sendBatchEntries(SqsClient amazonSQS, Exchange exchange, 
List<SendMessageBatchRequestEntry> entries) {
+        if (entries.isEmpty()) {
+            return;
+        }
+        int totalFailed = 0;
+        int totalSuccessful = 0;
+        SendMessageBatchResponse lastResult = null;
+        for (int i = 0; i < entries.size(); i += MAX_BATCH_ENTRIES) {
+            List<SendMessageBatchRequestEntry> chunk = entries.subList(i, 
Math.min(i + MAX_BATCH_ENTRIES, entries.size()));
+            SendMessageBatchRequest request = SendMessageBatchRequest.builder()
+                    .queueUrl(getQueueUrl())
+                    .entries(chunk)
+                    .build();
+            lastResult = amazonSQS.sendMessageBatch(request);
+            totalFailed += ObjectHelper.isNotEmpty(lastResult.failed()) ? 
lastResult.failed().size() : 0;
+            totalSuccessful += 
ObjectHelper.isNotEmpty(lastResult.successful()) ? 
lastResult.successful().size() : 0;
+        }
+        Message message = getMessageForResponse(exchange);
+        message.setBody(lastResult);
+        message.setHeader(Sqs2Constants.FAILED_MESSAGE_COUNT, totalFailed);
+        message.setHeader(Sqs2Constants.SUCCESSFUL_MESSAGE_COUNT, 
totalSuccessful);
+    }
+
     private void deleteMessage(SqsClient amazonSQS, Exchange exchange) {
         String receiptHandle = 
exchange.getIn().getHeader(Sqs2Constants.RECEIPT_HANDLE, String.class);
         DeleteMessageRequest.Builder request = DeleteMessageRequest.builder();

Reply via email to