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; }
