This is an automated email from the ASF dual-hosted git repository. davsclaus pushed a commit to branch fix/CAMEL-24154 in repository https://gitbox.apache.org/repos/asf/camel.git
commit d29abe8c9dccd4077d7509b62c3bec52c81f630c Author: Claus Ibsen <[email protected]> AuthorDate: Fri Jul 17 16:47:11 2026 +0200 CAMEL-24154: camel-aws2-ddb - Fix DDB Streams consumer silent data loss on resharding Co-Authored-By: Claude Opus 4.6 <[email protected]> Signed-off-by: Claus Ibsen <[email protected]> --- .../aws2/ddbstream/ShardIteratorHandler.java | 48 +++++++++- .../aws2/ddbstream/AmazonDDBStreamsClientMock.java | 9 ++ .../aws2/ddbstream/ShardIteratorHandlerTest.java | 105 +++++++++++++++++++++ 3 files changed, 157 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..e90a1330daa8 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,10 +16,13 @@ */ package org.apache.camel.component.aws2.ddbstream; +import java.util.ArrayList; import java.util.HashMap; +import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.Map.Entry; +import java.util.Set; import org.apache.camel.component.aws2.ddbstream.Ddb2StreamConfiguration.StreamIteratorType; import org.slf4j.Logger; @@ -38,6 +41,7 @@ class ShardIteratorHandler { private final Ddb2StreamEndpoint endpoint; private final ShardTree shardTree = new ShardTree(); + private final Set<String> pendingClosedShards = new HashSet<>(); private String streamArn; private Map<String, String> currentShardIterators = new HashMap<>(); @@ -52,13 +56,17 @@ 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()); + shardTree.populate(describeStreamPaginated()); StreamIteratorType streamIteratorType = getEndpoint().getConfiguration().getStreamIteratorType(); currentShardIterators = getCurrentShardIterators(streamIteratorType); } else { + // Refresh the shard tree when any tracked shard has closed, so that + // children created by DynamoDB's shard rotation become visible. + if (!pendingClosedShards.isEmpty()) { + shardTree.populate(describeStreamPaginated()); + } + Map<String, String> childShardIterators = new HashMap<>(); for (Entry<String, String> currentShardIterator : currentShardIterators.entrySet()) { List<Shard> children = shardTree.getChildren(currentShardIterator.getKey()); @@ -71,22 +79,37 @@ class ShardIteratorHandler { } } } + + for (String closedShardId : pendingClosedShards) { + for (Shard child : shardTree.getChildren(closedShardId)) { + if (!childShardIterators.containsKey(child.shardId())) { + String shardIterator = getShardIterator(child.shardId(), ShardIteratorType.TRIM_HORIZON); + childShardIterators.put(child.shardId(), shardIterator); + } + } + } + pendingClosedShards.clear(); + currentShardIterators = childShardIterators; } LOG.trace("Shard Iterators are: {}", currentShardIterators); - return currentShardIterators; + return new HashMap<>(currentShardIterators); } void updateShardIterator(String shardId, String nextShardIterator) { if (nextShardIterator == null) { // Shard has become inactive and all records have been consumed. currentShardIterators.remove(shardId); + pendingClosedShards.add(shardId); } else { currentShardIterators.put(shardId, nextShardIterator); } } String requestFreshShardIterator(String shardId, String lastSeenSequenceNumber) { - String shardIterator = getShardIterator(shardId, ShardIteratorType.AFTER_SEQUENCE_NUMBER, lastSeenSequenceNumber); + ShardIteratorType type = lastSeenSequenceNumber != null + ? ShardIteratorType.AFTER_SEQUENCE_NUMBER + : ShardIteratorType.TRIM_HORIZON; + String shardIterator = getShardIterator(shardId, type, lastSeenSequenceNumber); currentShardIterators.put(shardId, shardIterator); return shardIterator; } @@ -95,6 +118,21 @@ class ShardIteratorHandler { return endpoint; } + private List<Shard> describeStreamPaginated() { + List<Shard> allShards = new ArrayList<>(); + String lastEvaluatedShardId = null; + do { + DescribeStreamRequest.Builder builder = DescribeStreamRequest.builder().streamArn(streamArn); + if (lastEvaluatedShardId != null) { + builder.exclusiveStartShardId(lastEvaluatedShardId); + } + DescribeStreamResponse response = getClient().describeStream(builder.build()); + allShards.addAll(response.streamDescription().shards()); + lastEvaluatedShardId = response.streamDescription().lastEvaluatedShardId(); + } while (lastEvaluatedShardId != null); + return allShards; + } + private String getStreamArn() { ListStreamsResponse streamsListResult = getClient().listStreams( ListStreamsRequest.builder().tableName(getEndpoint().getConfiguration().getTableName()).build()); diff --git a/components/camel-aws/camel-aws2-ddb/src/test/java/org/apache/camel/component/aws2/ddbstream/AmazonDDBStreamsClientMock.java b/components/camel-aws/camel-aws2-ddb/src/test/java/org/apache/camel/component/aws2/ddbstream/AmazonDDBStreamsClientMock.java index 27f9e635b5f9..1fecf13656a5 100644 --- a/components/camel-aws/camel-aws2-ddb/src/test/java/org/apache/camel/component/aws2/ddbstream/AmazonDDBStreamsClientMock.java +++ b/components/camel-aws/camel-aws2-ddb/src/test/java/org/apache/camel/component/aws2/ddbstream/AmazonDDBStreamsClientMock.java @@ -16,7 +16,9 @@ */ package org.apache.camel.component.aws2.ddbstream; +import java.util.ArrayList; import java.util.HashMap; +import java.util.List; import java.util.Map; import java.util.Map.Entry; @@ -37,6 +39,7 @@ import static org.apache.camel.component.aws2.ddbstream.ShardFixtures.STREAM_ARN class AmazonDDBStreamsClientMock implements DynamoDbStreamsClient { private final Map<Shard, String> shardsToIterators = new HashMap<>(); + private final List<GetShardIteratorRequest> shardIteratorRequests = new ArrayList<>(); @Override public ListStreamsResponse listStreams(ListStreamsRequest listStreamsRequest) { @@ -56,6 +59,7 @@ class AmazonDDBStreamsClientMock implements DynamoDbStreamsClient { @Override public GetShardIteratorResponse getShardIterator(GetShardIteratorRequest request) { + shardIteratorRequests.add(request); String shardIterator = shardsToIterators.entrySet().stream() .filter(s -> s.getKey().shardId().equals(request.shardId())) .map(Entry::getValue) @@ -74,6 +78,11 @@ class AmazonDDBStreamsClientMock implements DynamoDbStreamsClient { } void setMockedShardAndIteratorResponse(Shard shard, String iterator) { + shardsToIterators.keySet().removeIf(s -> s.shardId().equals(shard.shardId())); shardsToIterators.put(shard, iterator); } + + List<GetShardIteratorRequest> getShardIteratorRequests() { + return shardIteratorRequests; + } } diff --git a/components/camel-aws/camel-aws2-ddb/src/test/java/org/apache/camel/component/aws2/ddbstream/ShardIteratorHandlerTest.java b/components/camel-aws/camel-aws2-ddb/src/test/java/org/apache/camel/component/aws2/ddbstream/ShardIteratorHandlerTest.java index d9451d582438..ef370fbc76eb 100644 --- a/components/camel-aws/camel-aws2-ddb/src/test/java/org/apache/camel/component/aws2/ddbstream/ShardIteratorHandlerTest.java +++ b/components/camel-aws/camel-aws2-ddb/src/test/java/org/apache/camel/component/aws2/ddbstream/ShardIteratorHandlerTest.java @@ -24,11 +24,15 @@ import org.apache.camel.component.aws2.ddbstream.Ddb2StreamConfiguration.StreamI import org.apache.camel.test.junit6.CamelTestSupport; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import software.amazon.awssdk.services.dynamodb.model.SequenceNumberRange; +import software.amazon.awssdk.services.dynamodb.model.Shard; +import software.amazon.awssdk.services.dynamodb.model.ShardIteratorType; import static org.apache.camel.component.aws2.ddbstream.ShardFixtures.*; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; class ShardIteratorHandlerTest extends CamelTestSupport { @@ -158,4 +162,105 @@ class ShardIteratorHandlerTest extends CamelTestSupport { assertThrows(IllegalArgumentException.class, () -> underTest.getShardIterators()); } + @Test + void shouldDiscoverNewShardsAfterResharding() throws Exception { + // Fresh mock with only two active shards + AmazonDDBStreamsClientMock reshardingMock = new AmazonDDBStreamsClientMock(); + component.getConfiguration().setAmazonDynamoDbStreamsClient(reshardingMock); + + Shard shardA = Shard.builder() + .shardId("SHARD_A") + .sequenceNumberRange(SequenceNumberRange.builder() + .startingSequenceNumber("100").build()) + .build(); + Shard shardB = Shard.builder() + .shardId("SHARD_B") + .sequenceNumberRange(SequenceNumberRange.builder() + .startingSequenceNumber("200").build()) + .build(); + String iterA = STREAM_ARN + "|iter-A"; + String iterB = STREAM_ARN + "|iter-B"; + + reshardingMock.setMockedShardAndIteratorResponse(shardA, iterA); + reshardingMock.setMockedShardAndIteratorResponse(shardB, iterB); + + component.getConfiguration().setStreamIteratorType(StreamIteratorType.FROM_LATEST); + Ddb2StreamEndpoint endpoint = (Ddb2StreamEndpoint) component.createEndpoint("aws2-ddbstreams://myTable"); + ShardIteratorHandler underTest = new ShardIteratorHandler(endpoint); + endpoint.doStart(); + + // Initial poll: both leaves are returned + Map<String, String> iter1 = underTest.getShardIterators(); + assertEquals(2, iter1.size()); + assertTrue(iter1.containsKey("SHARD_A")); + assertTrue(iter1.containsKey("SHARD_B")); + + // Shard A closes (consumer drained it) + underTest.updateShardIterator("SHARD_A", null); + + // Simulate resharding: A becomes closed, children A1 and A2 appear + Shard closedA = Shard.builder() + .shardId("SHARD_A") + .sequenceNumberRange(SequenceNumberRange.builder() + .startingSequenceNumber("100").endingSequenceNumber("150").build()) + .build(); + Shard shardA1 = Shard.builder() + .shardId("SHARD_A1") + .parentShardId("SHARD_A") + .sequenceNumberRange(SequenceNumberRange.builder() + .startingSequenceNumber("151").build()) + .build(); + Shard shardA2 = Shard.builder() + .shardId("SHARD_A2") + .parentShardId("SHARD_A") + .sequenceNumberRange(SequenceNumberRange.builder() + .startingSequenceNumber("152").build()) + .build(); + String iterA1 = STREAM_ARN + "|iter-A1"; + String iterA2 = STREAM_ARN + "|iter-A2"; + reshardingMock.setMockedShardAndIteratorResponse(closedA, iterA); + reshardingMock.setMockedShardAndIteratorResponse(shardA1, iterA1); + reshardingMock.setMockedShardAndIteratorResponse(shardA2, iterA2); + + // Next poll: tree is refreshed, children of A are discovered + Map<String, String> iter2 = underTest.getShardIterators(); + assertEquals(3, iter2.size()); + assertTrue(iter2.containsKey("SHARD_A1")); + assertTrue(iter2.containsKey("SHARD_A2")); + assertTrue(iter2.containsKey("SHARD_B")); + assertFalse(iter2.containsKey("SHARD_A")); + } + + @Test + void shouldReturnDefensiveCopyOfShardIterators() throws Exception { + component.getConfiguration().setStreamIteratorType(StreamIteratorType.FROM_LATEST); + Ddb2StreamEndpoint endpoint = (Ddb2StreamEndpoint) component.createEndpoint("aws2-ddbstreams://myTable"); + ShardIteratorHandler underTest = new ShardIteratorHandler(endpoint); + endpoint.doStart(); + + Map<String, String> first = underTest.getShardIterators(); + Map<String, String> second = underTest.getShardIterators(); + + // Mutating the returned map must not affect internal state + first.put("EXTRA_SHARD", "EXTRA_ITERATOR"); + assertFalse(second.containsKey("EXTRA_SHARD")); + } + + @Test + void shouldUseTrimHorizonWhenSequenceNumberIsNull() throws Exception { + component.getConfiguration().setStreamIteratorType(StreamIteratorType.FROM_LATEST); + Ddb2StreamEndpoint endpoint = (Ddb2StreamEndpoint) component.createEndpoint("aws2-ddbstreams://myTable"); + ShardIteratorHandler underTest = new ShardIteratorHandler(endpoint); + endpoint.doStart(); + + underTest.getShardIterators(); + dynamoDbStreamsClient.getShardIteratorRequests().clear(); + + underTest.requestFreshShardIterator(SHARD_3.shardId(), null); + + assertEquals(1, dynamoDbStreamsClient.getShardIteratorRequests().size()); + assertEquals(ShardIteratorType.TRIM_HORIZON, + dynamoDbStreamsClient.getShardIteratorRequests().get(0).shardIteratorType()); + } + }
