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 c874ecbb1e19 CAMEL-24156: Fix high-severity findings across camel-aws2 
components
c874ecbb1e19 is described below

commit c874ecbb1e19761c9aa3daa1ba97b4bc6b98a60a
Author: Claus Ibsen <[email protected]>
AuthorDate: Sat Jul 18 11:39:38 2026 +0200

    CAMEL-24156: Fix high-severity findings across camel-aws2 components
    
    - S3: Close ResponseInputStream when ignoreBody=true to prevent connection
      pool exhaustion; use dynamic bucket name in all multipart upload API 
calls;
      close FileInputStream on contentLength < partSize shortcut path
    - SQS: Skip per-message DelaySeconds when delayQueue=true or queue is FIFO
    - SNS: Set ContentBasedDeduplication on auto-created FIFO topics
    - Kinesis: Initialize shardClosed to documented default (ignore); handle
      ExpiredIteratorException with fresh iterator retry; reset client fields
      in KinesisConnection.close(); create per-shard KinesisResumeAction via
      Camel Injector API preserving registry-based customization
    - DDB: Use Number instead of Integer/Double checks in JSON transformer so
      Long, BigDecimal, Float values are stored as DynamoDB N type
    
    Closes #24863
    Co-Authored-By: Claude Opus 4.6 <[email protected]>
---
 .../ddb/transform/Ddb2JsonDataTypeTransformer.java |  6 +-
 .../src/main/docs/aws2-kinesis-component.adoc      | 11 ++++
 .../aws2/kinesis/Kinesis2Configuration.java        |  2 +-
 .../component/aws2/kinesis/Kinesis2Consumer.java   | 69 +++++++++++++++-------
 .../component/aws2/kinesis/KinesisConnection.java  | 17 ++++--
 .../aws2/kinesis/consumer/KinesisResumeAction.java |  3 +
 .../integration/KinesisConsumerResumeIT.java       | 22 +++----
 .../camel/component/aws2/s3/AWS2S3Producer.java    | 22 ++++---
 .../camel/component/aws2/sns/Sns2Endpoint.java     |  3 +
 .../camel/component/aws2/sqs/Sqs2Producer.java     |  6 ++
 .../ROOT/pages/camel-4x-upgrade-guide-4_22.adoc    | 17 ++++++
 11 files changed, 128 insertions(+), 50 deletions(-)

diff --git 
a/components/camel-aws/camel-aws2-ddb/src/main/java/org/apache/camel/component/aws2/ddb/transform/Ddb2JsonDataTypeTransformer.java
 
b/components/camel-aws/camel-aws2-ddb/src/main/java/org/apache/camel/component/aws2/ddb/transform/Ddb2JsonDataTypeTransformer.java
index 80f2e2aac2ee..99a16d878d3c 100644
--- 
a/components/camel-aws/camel-aws2-ddb/src/main/java/org/apache/camel/component/aws2/ddb/transform/Ddb2JsonDataTypeTransformer.java
+++ 
b/components/camel-aws/camel-aws2-ddb/src/main/java/org/apache/camel/component/aws2/ddb/transform/Ddb2JsonDataTypeTransformer.java
@@ -194,11 +194,7 @@ public class Ddb2JsonDataTypeTransformer extends 
Transformer {
             return AttributeValue.builder().s(value.toString()).build();
         }
 
-        if (value instanceof Integer) {
-            return AttributeValue.builder().n(value.toString()).build();
-        }
-
-        if (value instanceof Double) {
+        if (value instanceof Number) {
             return AttributeValue.builder().n(value.toString()).build();
         }
 
diff --git 
a/components/camel-aws/camel-aws2-kinesis/src/main/docs/aws2-kinesis-component.adoc
 
b/components/camel-aws/camel-aws2-kinesis/src/main/docs/aws2-kinesis-component.adoc
index d6ee55c5d925..d5b7fd9e67cc 100644
--- 
a/components/camel-aws/camel-aws2-kinesis/src/main/docs/aws2-kinesis-component.adoc
+++ 
b/components/camel-aws/camel-aws2-kinesis/src/main/docs/aws2-kinesis-component.adoc
@@ -66,6 +66,17 @@ all available shards (multiple shards consumption) of Amazon 
Kinesis, therefore,
 property in the DSL configuration empty, then it'll consume all available 
shards
 otherwise only the specified shard corresponding to the shardId will be 
consumed.
 
+=== Custom Resume Action
+
+The consumer supports a custom `KinesisResumeAction` for controlling where 
each shard starts reading
+(e.g., resuming from a persisted sequence number). Register a subclass in the 
Camel registry under the
+key `CamelKinesisDbResumeAction`.
+
+The consumer creates a separate instance per shard via reflection to avoid 
concurrent mutation when
+multiple shards are processed in parallel. Because of this, custom subclasses 
*must* provide a
+public no-arg constructor. Per-shard state (`builder`, `shardId`, 
`streamName`) is injected via
+setters after construction.
+
 === Batch Producer
 
 This component implements the Batch Producer.
diff --git 
a/components/camel-aws/camel-aws2-kinesis/src/main/java/org/apache/camel/component/aws2/kinesis/Kinesis2Configuration.java
 
b/components/camel-aws/camel-aws2-kinesis/src/main/java/org/apache/camel/component/aws2/kinesis/Kinesis2Configuration.java
index 4b214f845fec..fe6942707717 100644
--- 
a/components/camel-aws/camel-aws2-kinesis/src/main/java/org/apache/camel/component/aws2/kinesis/Kinesis2Configuration.java
+++ 
b/components/camel-aws/camel-aws2-kinesis/src/main/java/org/apache/camel/component/aws2/kinesis/Kinesis2Configuration.java
@@ -68,7 +68,7 @@ public class Kinesis2Configuration implements Cloneable, 
AwsCommonConfiguration
                             + " In case of ignore a WARN message will be 
logged once and the consumer will not process new messages until restarted,"
                             + "in case of silent there will be no logging and 
the consumer will not process new messages until restarted,"
                             + "in case of fail a ReachedClosedStateException 
will be thrown")
-    private Kinesis2ShardClosedStrategyEnum shardClosed;
+    private Kinesis2ShardClosedStrategyEnum shardClosed = 
Kinesis2ShardClosedStrategyEnum.ignore;
     @UriParam(label = "proxy", enums = "HTTP,HTTPS", defaultValue = "HTTPS",
               description = "To define a proxy protocol when instantiating the 
Kinesis client")
     private Protocol proxyProtocol = Protocol.HTTPS;
diff --git 
a/components/camel-aws/camel-aws2-kinesis/src/main/java/org/apache/camel/component/aws2/kinesis/Kinesis2Consumer.java
 
b/components/camel-aws/camel-aws2-kinesis/src/main/java/org/apache/camel/component/aws2/kinesis/Kinesis2Consumer.java
index c8ce34fa8dd3..f85bcb7de6ce 100644
--- 
a/components/camel-aws/camel-aws2-kinesis/src/main/java/org/apache/camel/component/aws2/kinesis/Kinesis2Consumer.java
+++ 
b/components/camel-aws/camel-aws2-kinesis/src/main/java/org/apache/camel/component/aws2/kinesis/Kinesis2Consumer.java
@@ -43,6 +43,7 @@ import org.apache.camel.util.CastUtils;
 import org.apache.camel.util.ObjectHelper;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
+import software.amazon.awssdk.services.kinesis.model.ExpiredIteratorException;
 import software.amazon.awssdk.services.kinesis.model.GetRecordsRequest;
 import software.amazon.awssdk.services.kinesis.model.GetRecordsResponse;
 import software.amazon.awssdk.services.kinesis.model.GetShardIteratorRequest;
@@ -152,6 +153,42 @@ public class Kinesis2Consumer extends 
ScheduledBatchPollingConsumer implements R
             return;
         }
 
+        GetRecordsResponse result;
+        try {
+            result = getRecords(shardIterator, kinesisConnection);
+        } catch (ExpiredIteratorException e) {
+            LOG.warn("Shard iterator expired for shard {} on stream {}, 
requesting a fresh one",
+                    shard.shardId(), 
getEndpoint().getConfiguration().getStreamName());
+            currentShardIterators.remove(shard.shardId());
+            try {
+                shardIterator = getShardIterator(shard, kinesisConnection);
+            } catch (InterruptedException ie) {
+                Thread.currentThread().interrupt();
+                throw new RuntimeException(ie);
+            } catch (ExecutionException ee) {
+                throw new RuntimeException(ee);
+            }
+            if (ObjectHelper.isEmpty(shardIterator)) {
+                return;
+            }
+            result = getRecords(shardIterator, kinesisConnection);
+        }
+
+        try {
+            Queue<Exchange> exchanges = createExchanges(shard, 
result.records());
+            
processedExchangeCount.addAndGet(processBatch(CastUtils.cast(exchanges)));
+        } catch (Exception e) {
+            throw new RuntimeException(e);
+        }
+
+        // May cache the last successful sequence number, and pass it to the
+        // getRecords request. That way, on the next poll, we start from where
+        // we left off, however, I don't know what happens to subsequent
+        // exchanges when an earlier exchange fails.
+        updateShardIterator(shard, result.nextShardIterator());
+    }
+
+    private GetRecordsResponse getRecords(String shardIterator, 
KinesisConnection kinesisConnection) {
         GetRecordsRequest req = GetRecordsRequest
                 .builder()
                 .shardIterator(shardIterator)
@@ -160,10 +197,9 @@ public class Kinesis2Consumer extends 
ScheduledBatchPollingConsumer implements R
                         .getMaxResultsPerRequest())
                 .build();
 
-        GetRecordsResponse result;
         if (getEndpoint().getConfiguration().isAsyncClient()) {
             try {
-                result = kinesisConnection
+                return kinesisConnection
                         .getAsyncClient(getEndpoint())
                         .getRecords(req)
                         .get();
@@ -171,26 +207,16 @@ public class Kinesis2Consumer extends 
ScheduledBatchPollingConsumer implements R
                 Thread.currentThread().interrupt();
                 throw new RuntimeException(e);
             } catch (ExecutionException e) {
+                if (e.getCause() instanceof ExpiredIteratorException ex) {
+                    throw ex;
+                }
                 throw new RuntimeException(e);
             }
         } else {
-            result = kinesisConnection
+            return kinesisConnection
                     .getClient(getEndpoint())
                     .getRecords(req);
         }
-
-        try {
-            Queue<Exchange> exchanges = createExchanges(shard, 
result.records());
-            
processedExchangeCount.addAndGet(processBatch(CastUtils.cast(exchanges)));
-        } catch (Exception e) {
-            throw new RuntimeException(e);
-        }
-
-        // May cache the last successful sequence number, and pass it to the
-        // getRecords request. That way, on the next poll, we start from where
-        // we left off, however, I don't know what happens to subsequent
-        // exchanges when an earlier exchange fails.
-        updateShardIterator(shard, result.nextShardIterator());
     }
 
     private void updateShardIterator(Shard shard, String nextShardIterator) {
@@ -315,15 +341,16 @@ public class Kinesis2Consumer extends 
ScheduledBatchPollingConsumer implements R
     }
 
     private KinesisResumeAction resolveResumeAction(String shardId, 
GetShardIteratorRequest.Builder req) {
-        KinesisResumeAction action
-                = 
getEndpoint().getCamelContext().getRegistry().lookupByNameAndType(Kinesis2Constants.RESUME_ACTION,
-                        KinesisResumeAction.class);
-        if (ObjectHelper.isEmpty(action)) {
+        KinesisResumeAction template = 
getEndpoint().getCamelContext().getRegistry()
+                .lookupByNameAndType(Kinesis2Constants.RESUME_ACTION, 
KinesisResumeAction.class);
+        KinesisResumeAction action;
+        if (ObjectHelper.isEmpty(template)) {
             action = new KinesisResumeAction(req);
         } else {
+            action = (KinesisResumeAction) 
getEndpoint().getCamelContext().getInjector()
+                    .newInstance(template.getClass());
             action.setBuilder(req);
         }
-
         action.setShardId(shardId);
         action.setStreamName(getEndpoint().getConfiguration().getStreamName());
         return action;
diff --git 
a/components/camel-aws/camel-aws2-kinesis/src/main/java/org/apache/camel/component/aws2/kinesis/KinesisConnection.java
 
b/components/camel-aws/camel-aws2-kinesis/src/main/java/org/apache/camel/component/aws2/kinesis/KinesisConnection.java
index b9e888f3728e..64d32eaf34a1 100644
--- 
a/components/camel-aws/camel-aws2-kinesis/src/main/java/org/apache/camel/component/aws2/kinesis/KinesisConnection.java
+++ 
b/components/camel-aws/camel-aws2-kinesis/src/main/java/org/apache/camel/component/aws2/kinesis/KinesisConnection.java
@@ -73,11 +73,18 @@ public class KinesisConnection implements Closeable {
 
     @Override
     public void close() throws IOException {
-        if (ObjectHelper.isNotEmpty(kinesisClient)) {
-            kinesisClient.close();
-        }
-        if (ObjectHelper.isNotEmpty(kinesisAsyncClient)) {
-            kinesisAsyncClient.close();
+        lock.lock();
+        try {
+            if (ObjectHelper.isNotEmpty(kinesisClient)) {
+                kinesisClient.close();
+                kinesisClient = null;
+            }
+            if (ObjectHelper.isNotEmpty(kinesisAsyncClient)) {
+                kinesisAsyncClient.close();
+                kinesisAsyncClient = null;
+            }
+        } finally {
+            lock.unlock();
         }
     }
 }
diff --git 
a/components/camel-aws/camel-aws2-kinesis/src/main/java/org/apache/camel/component/aws2/kinesis/consumer/KinesisResumeAction.java
 
b/components/camel-aws/camel-aws2-kinesis/src/main/java/org/apache/camel/component/aws2/kinesis/consumer/KinesisResumeAction.java
index 11dcab57ae7f..5501dbd17e72 100644
--- 
a/components/camel-aws/camel-aws2-kinesis/src/main/java/org/apache/camel/component/aws2/kinesis/consumer/KinesisResumeAction.java
+++ 
b/components/camel-aws/camel-aws2-kinesis/src/main/java/org/apache/camel/component/aws2/kinesis/consumer/KinesisResumeAction.java
@@ -21,6 +21,9 @@ import org.apache.camel.resume.ResumeAction;
 import software.amazon.awssdk.services.kinesis.model.GetShardIteratorRequest;
 import software.amazon.awssdk.services.kinesis.model.ShardIteratorType;
 
+/**
+ * Subclasses must provide a public no-arg constructor — the consumer creates 
per-shard instances via reflection.
+ */
 public class KinesisResumeAction implements ResumeAction {
 
     private GetShardIteratorRequest.Builder builder;
diff --git 
a/components/camel-aws/camel-aws2-kinesis/src/test/java/org/apache/camel/component/aws2/kinesis/integration/KinesisConsumerResumeIT.java
 
b/components/camel-aws/camel-aws2-kinesis/src/test/java/org/apache/camel/component/aws2/kinesis/integration/KinesisConsumerResumeIT.java
index 18cf1fa1c0d2..eb43ec446ec8 100644
--- 
a/components/camel-aws/camel-aws2-kinesis/src/test/java/org/apache/camel/component/aws2/kinesis/integration/KinesisConsumerResumeIT.java
+++ 
b/components/camel-aws/camel-aws2-kinesis/src/test/java/org/apache/camel/component/aws2/kinesis/integration/KinesisConsumerResumeIT.java
@@ -77,20 +77,19 @@ public class KinesisConsumerResumeIT extends 
CamelTestSupport {
         }
     }
 
-    private static final class TestResumeAction extends KinesisResumeAction {
-        private List<PutRecordsResponse> previousRecords;
-        private final int expectedCount;
+    public static class TestResumeAction extends KinesisResumeAction {
+        private static volatile List<PutRecordsResponse> previousRecords;
+        private static volatile int expectedCount;
 
-        private TestResumeAction(int expectedCount) {
-            this.expectedCount = expectedCount;
+        public TestResumeAction() {
         }
 
-        public void setPreviousRecords(List<PutRecordsResponse> 
previousRecords) {
-            this.previousRecords = previousRecords;
+        public static void setPreviousRecords(List<PutRecordsResponse> 
records) {
+            previousRecords = records;
         }
 
-        public int getExpectedCount() {
-            return expectedCount;
+        public static void setExpectedCount(int count) {
+            expectedCount = count;
         }
 
         @Override
@@ -125,7 +124,7 @@ public class KinesisConsumerResumeIT extends 
CamelTestSupport {
     private final int expectedCount = messageCount / 2;
     private List<KinesisData> receivedMessages = new CopyOnWriteArrayList<>();
     private List<PutRecordsResponse> previousRecords;
-    private TestResumeAction action = new TestResumeAction(expectedCount);
+    private TestResumeAction action = new TestResumeAction();
 
     @Override
     protected RouteBuilder createRouteBuilder() {
@@ -176,7 +175,8 @@ public class KinesisConsumerResumeIT extends 
CamelTestSupport {
             }
         }
 
-        action.setPreviousRecords(previousRecords);
+        TestResumeAction.setExpectedCount(expectedCount);
+        TestResumeAction.setPreviousRecords(previousRecords);
     }
 
     @AfterEach
diff --git 
a/components/camel-aws/camel-aws2-s3/src/main/java/org/apache/camel/component/aws2/s3/AWS2S3Producer.java
 
b/components/camel-aws/camel-aws2-s3/src/main/java/org/apache/camel/component/aws2/s3/AWS2S3Producer.java
index eecb6949cf74..d0ff4e1800fc 100644
--- 
a/components/camel-aws/camel-aws2-s3/src/main/java/org/apache/camel/component/aws2/s3/AWS2S3Producer.java
+++ 
b/components/camel-aws/camel-aws2-s3/src/main/java/org/apache/camel/component/aws2/s3/AWS2S3Producer.java
@@ -214,7 +214,11 @@ public class AWS2S3Producer extends DefaultProducer {
         if (contentLength == 0 || contentLength < partSize) {
             // optimize to do a single op if content length is known and < 
part size
             LOG.debug("Payload size < partSize ({} > {}). Uploading payload in 
single operation", contentLength, partSize);
-            doPutObject(exchange, objectMetadata, filePayload, inputStream, 
contentLength);
+            try {
+                doPutObject(exchange, objectMetadata, filePayload, 
inputStream, contentLength);
+            } finally {
+                IOHelper.close(inputStream);
+            }
             return;
         }
 
@@ -224,7 +228,7 @@ public class AWS2S3Producer extends DefaultProducer {
         final String keyName = AWS2S3Utils.determineKey(exchange, 
getConfiguration());
         final String bucketName = AWS2S3Utils.determineBucketName(exchange, 
getConfiguration());
         CreateMultipartUploadRequest.Builder createMultipartUploadRequest
-                = 
CreateMultipartUploadRequest.builder().bucket(getConfiguration().getBucketName()).key(keyName);
+                = 
CreateMultipartUploadRequest.builder().bucket(bucketName).key(keyName);
 
         String storageClass = AWS2S3Utils.determineStorageClass(exchange, 
getConfiguration());
         if (ObjectHelper.isNotEmpty(storageClass)) {
@@ -279,7 +283,7 @@ public class AWS2S3Producer extends DefaultProducer {
             for (int part = 1; position < contentLength; part++) {
                 partSize = Math.min(partSize, contentLength - position);
 
-                UploadPartRequest uploadRequest = 
UploadPartRequest.builder().bucket(getConfiguration().getBucketName())
+                UploadPartRequest uploadRequest = 
UploadPartRequest.builder().bucket(bucketName)
                         .key(keyName).uploadId(initResponse.uploadId())
                         .partNumber(part).build();
 
@@ -295,7 +299,7 @@ public class AWS2S3Producer extends DefaultProducer {
             LOG.debug("Completing multi-part upload for {}", keyName);
             CompletedMultipartUpload completeMultipartUpload = 
CompletedMultipartUpload.builder().parts(completedParts).build();
             CompleteMultipartUploadRequest.Builder compRequestBuilder = 
CompleteMultipartUploadRequest.builder()
-                    
.multipartUpload(completeMultipartUpload).bucket(getConfiguration().getBucketName()).key(keyName)
+                    
.multipartUpload(completeMultipartUpload).bucket(bucketName).key(keyName)
                     .uploadId(initResponse.uploadId());
             if (getConfiguration().isConditionalWritesEnabled()) {
                 compRequestBuilder.ifNoneMatch("*");
@@ -304,7 +308,7 @@ public class AWS2S3Producer extends DefaultProducer {
 
         } catch (Exception e) {
             getEndpoint().getS3Client()
-                    
.abortMultipartUpload(AbortMultipartUploadRequest.builder().bucket(getConfiguration().getBucketName())
+                    
.abortMultipartUpload(AbortMultipartUploadRequest.builder().bucket(bucketName)
                             
.key(keyName).uploadId(initResponse.uploadId()).build());
             throw e;
         } finally {
@@ -616,7 +620,9 @@ public class AWS2S3Producer extends DefaultProducer {
                 ResponseInputStream<GetObjectResponse> res
                         = s3Client.getObject(req, 
ResponseTransformer.toInputStream());
                 Message message = getMessageForResponse(exchange);
-                if (!getConfiguration().isIgnoreBody()) {
+                if (getConfiguration().isIgnoreBody()) {
+                    IOHelper.close(res);
+                } else {
                     message.setBody(res);
                 }
                 populateMetadata(res, message);
@@ -658,7 +664,9 @@ public class AWS2S3Producer extends DefaultProducer {
             ResponseInputStream<GetObjectResponse> res = 
s3Client.getObject(req.build(), ResponseTransformer.toInputStream());
 
             Message message = getMessageForResponse(exchange);
-            if (!getConfiguration().isIgnoreBody()) {
+            if (getConfiguration().isIgnoreBody()) {
+                IOHelper.close(res);
+            } else {
                 message.setBody(res);
             }
             populateMetadata(res, message);
diff --git 
a/components/camel-aws/camel-aws2-sns/src/main/java/org/apache/camel/component/aws2/sns/Sns2Endpoint.java
 
b/components/camel-aws/camel-aws2-sns/src/main/java/org/apache/camel/component/aws2/sns/Sns2Endpoint.java
index 3c64b5fa5c10..e5cdb788e7d5 100644
--- 
a/components/camel-aws/camel-aws2-sns/src/main/java/org/apache/camel/component/aws2/sns/Sns2Endpoint.java
+++ 
b/components/camel-aws/camel-aws2-sns/src/main/java/org/apache/camel/component/aws2/sns/Sns2Endpoint.java
@@ -149,6 +149,9 @@ public class Sns2Endpoint extends DefaultEndpoint 
implements HeaderFilterStrateg
 
             if (configuration.isFifoTopic()) {
                 attributes.put("FifoTopic", "true");
+                if (configuration.getMessageDeduplicationIdStrategy() 
instanceof NullMessageDeduplicationIdStrategy) {
+                    attributes.put("ContentBasedDeduplication", "true");
+                }
                 builder.attributes(attributes);
             }
 
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 a911d04871bb..2dfeb90d2f0f 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
@@ -281,6 +281,9 @@ public class Sqs2Producer extends DefaultProducer {
     }
 
     private void addDelay(SendMessageRequest.Builder request, Exchange 
exchange) {
+        if (getEndpoint().getConfiguration().isDelayQueue() || 
getEndpoint().getConfiguration().isFifoQueue()) {
+            return;
+        }
         Integer headerValue = 
exchange.getIn().getHeader(Sqs2Constants.DELAY_HEADER, Integer.class);
         Integer delayValue;
         if (ObjectHelper.isEmpty(headerValue)) {
@@ -297,6 +300,9 @@ public class Sqs2Producer extends DefaultProducer {
     }
 
     private void addDelay(SendMessageBatchRequestEntry.Builder request, 
Exchange exchange) {
+        if (getEndpoint().getConfiguration().isDelayQueue() || 
getEndpoint().getConfiguration().isFifoQueue()) {
+            return;
+        }
         Integer headerValue = 
exchange.getIn().getHeader(Sqs2Constants.DELAY_HEADER, Integer.class);
         Integer delayValue;
         if (ObjectHelper.isEmpty(headerValue)) {
diff --git 
a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_22.adoc 
b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_22.adoc
index 3b2c9d782943..adf398e28d3a 100644
--- a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_22.adoc
+++ b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_22.adoc
@@ -623,6 +623,23 @@ Two behavior changes follow:
 * On startup the consumer tails events from the moment it starts, rather than 
first replaying the
   single most-recent historical event.
 
+=== camel-aws2-kinesis - Custom KinesisResumeAction requires a public no-arg 
constructor
+
+Custom `KinesisResumeAction` subclasses registered in the Camel registry under 
the
+`CamelKinesisDbResumeAction` key must now provide a public no-arg constructor. 
The consumer creates
+a separate instance per shard via reflection to avoid concurrent mutation when 
multiple shards are
+processed in parallel. Per-shard state (`builder`, `shardId`, `streamName`) is 
injected via setters
+after construction. Any initialization that was previously done in a 
parameterized constructor should
+be moved to setters or to the `evalEntry` method.
+
+=== camel-aws2-sqs - per-message delay skipped for delay queues and FIFO queues
+
+The SQS producer no longer sets per-message `DelaySeconds` when 
`delayQueue=true` (delay is queue-level)
+or when the queue is FIFO (AWS rejects per-message delay with 
`InvalidParameterValue`). Previously the
+`CamelAwsSqsDelayHeader` header was applied unconditionally, which could cause 
AWS errors on FIFO queues
+and was redundant on delay queues. If your route relied on per-message delay 
overrides on a delay queue,
+the override will now be silently ignored.
+
 === camel-xslt-saxon - secure processing now applied unconditionally
 
 The `secureProcessing` option (default `true`) is now applied to the Saxon 
`TransformerFactory`

Reply via email to