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 56cd8dfdbf4b CAMEL-24154: camel-aws2-ddb - Fix DDB Streams consumer
silent data loss on resharding
56cd8dfdbf4b is described below
commit 56cd8dfdbf4be6d3a0df70ff169dca704fca6541
Author: Claus Ibsen <[email protected]>
AuthorDate: Fri Jul 17 17:27:23 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());
+ }
+
}