This is an automated email from the ASF dual-hosted git repository.

davsclaus pushed a commit to branch fix/CAMEL-24156
in repository https://gitbox.apache.org/repos/asf/camel.git

commit 52cb659e8f752a558b1d32329656a2f5cacaba90
Author: Claus Ibsen <[email protected]>
AuthorDate: Fri Jul 17 17:22:19 2026 +0200

    CAMEL-24156: camel-aws2-ddb - Fix stream consumer CME, pagination, and 
expired iterator
    
    - Return a defensive copy from getShardIterators() to prevent
      ConcurrentModificationException when updateShardIterator modifies the
      map during poll iteration with multiple shards
    - Paginate DescribeStream responses using lastEvaluatedShardId so streams
      with >100 shards are fully discovered
    - Fall back to TRIM_HORIZON when recovering from an expired iterator with
      no previously seen sequence number (null), instead of sending
      AFTER_SEQUENCE_NUMBER with null which causes ValidationException
    
    Co-Authored-By: Claude Opus 4.6 <[email protected]>
    Signed-off-by: Claus Ibsen <[email protected]>
---
 .../aws2/ddbstream/ShardIteratorHandler.java       | 25 +++++++++++++++++-----
 1 file changed, 20 insertions(+), 5 deletions(-)

diff --git 
a/components/camel-aws/camel-aws2-ddb/src/main/java/org/apache/camel/component/aws2/ddbstream/ShardIteratorHandler.java
 
b/components/camel-aws/camel-aws2-ddb/src/main/java/org/apache/camel/component/aws2/ddbstream/ShardIteratorHandler.java
index 685ea53d69fa..4add80151f45 100644
--- 
a/components/camel-aws/camel-aws2-ddb/src/main/java/org/apache/camel/component/aws2/ddbstream/ShardIteratorHandler.java
+++ 
b/components/camel-aws/camel-aws2-ddb/src/main/java/org/apache/camel/component/aws2/ddbstream/ShardIteratorHandler.java
@@ -16,6 +16,7 @@
  */
 package org.apache.camel.component.aws2.ddbstream;
 
+import java.util.ArrayList;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
@@ -52,9 +53,18 @@ class ShardIteratorHandler {
         }
         // Either return cached ones or get new ones via GetShardIterator 
requests.
         if (currentShardIterators.isEmpty()) {
-            DescribeStreamResponse streamDescriptionResult
-                    = 
getClient().describeStream(DescribeStreamRequest.builder().streamArn(streamArn).build());
-            
shardTree.populate(streamDescriptionResult.streamDescription().shards());
+            List<Shard> allShards = new ArrayList<>();
+            String exclusiveStartShardId = null;
+            do {
+                DescribeStreamRequest.Builder describeRequest = 
DescribeStreamRequest.builder().streamArn(streamArn);
+                if (exclusiveStartShardId != null) {
+                    
describeRequest.exclusiveStartShardId(exclusiveStartShardId);
+                }
+                DescribeStreamResponse streamDescriptionResult = 
getClient().describeStream(describeRequest.build());
+                
allShards.addAll(streamDescriptionResult.streamDescription().shards());
+                exclusiveStartShardId = 
streamDescriptionResult.streamDescription().lastEvaluatedShardId();
+            } while (exclusiveStartShardId != null);
+            shardTree.populate(allShards);
 
             StreamIteratorType streamIteratorType = 
getEndpoint().getConfiguration().getStreamIteratorType();
             currentShardIterators = 
getCurrentShardIterators(streamIteratorType);
@@ -74,7 +84,7 @@ class ShardIteratorHandler {
             currentShardIterators = childShardIterators;
         }
         LOG.trace("Shard Iterators are: {}", currentShardIterators);
-        return currentShardIterators;
+        return new HashMap<>(currentShardIterators);
     }
 
     void updateShardIterator(String shardId, String nextShardIterator) {
@@ -86,7 +96,12 @@ class ShardIteratorHandler {
     }
 
     String requestFreshShardIterator(String shardId, String 
lastSeenSequenceNumber) {
-        String shardIterator = getShardIterator(shardId, 
ShardIteratorType.AFTER_SEQUENCE_NUMBER, lastSeenSequenceNumber);
+        String shardIterator;
+        if (lastSeenSequenceNumber != null) {
+            shardIterator = getShardIterator(shardId, 
ShardIteratorType.AFTER_SEQUENCE_NUMBER, lastSeenSequenceNumber);
+        } else {
+            shardIterator = getShardIterator(shardId, 
ShardIteratorType.TRIM_HORIZON);
+        }
         currentShardIterators.put(shardId, shardIterator);
         return shardIterator;
     }

Reply via email to