This is an automated email from the ASF dual-hosted git repository.
ferenc-csaky pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/flink-connector-aws.git
The following commit(s) were added to refs/heads/main by this push:
new 6787a1d [FLINK-40184] Remove `DescribeStream`
inconsistency-resolution check from shard discovery
6787a1d is described below
commit 6787a1d91713423b15e3170ca3b54692d5a6ec95
Author: riyarawat-amazon <[email protected]>
AuthorDate: Thu Jul 23 00:27:53 2026 +0530
[FLINK-40184] Remove `DescribeStream` inconsistency-resolution check from
shard discovery
---
.../DynamodbStreamsSourceConfigConstants.java | 7 --
.../DynamoDbStreamsSourceEnumerator.java | 64 +-----------
.../tracker/SplitGraphInconsistencyTracker.java | 116 ---------------------
.../dynamodb/source/util/ListShardsResult.java | 17 +--
.../connector/dynamodb/source/util/ShardUtils.java | 23 ----
.../SplitGraphInconsistencyTrackerTest.java | 96 -----------------
.../dynamodb/source/util/ListShardsResultTest.java | 7 --
.../dynamodb/source/util/ShardUtilsTest.java | 9 --
8 files changed, 3 insertions(+), 336 deletions(-)
diff --git
a/flink-connector-aws/flink-connector-dynamodb/src/main/java/org/apache/flink/connector/dynamodb/source/config/DynamodbStreamsSourceConfigConstants.java
b/flink-connector-aws/flink-connector-dynamodb/src/main/java/org/apache/flink/connector/dynamodb/source/config/DynamodbStreamsSourceConfigConstants.java
index 8e21d09..94b02b8 100644
---
a/flink-connector-aws/flink-connector-dynamodb/src/main/java/org/apache/flink/connector/dynamodb/source/config/DynamodbStreamsSourceConfigConstants.java
+++
b/flink-connector-aws/flink-connector-dynamodb/src/main/java/org/apache/flink/connector/dynamodb/source/config/DynamodbStreamsSourceConfigConstants.java
@@ -45,13 +45,6 @@ public class DynamodbStreamsSourceConfigConstants {
.defaultValue(Duration.ofSeconds(60))
.withDescription("The interval between each attempt to
discover new shards.");
- public static final ConfigOption<Integer>
DESCRIBE_STREAM_INCONSISTENCY_RESOLUTION_RETRY_COUNT =
-
ConfigOptions.key("flink.describestream.inconsistencyresolution.retries")
- .intType()
- .defaultValue(5)
- .withDescription(
- "The number of times to retry build shard lineage
if describestream returns inconsistent response");
-
public static final ConfigOption<Integer> DYNAMODB_STREAMS_RETRY_COUNT =
ConfigOptions.key("flink.dynamodbstreams.numretries")
.intType()
diff --git
a/flink-connector-aws/flink-connector-dynamodb/src/main/java/org/apache/flink/connector/dynamodb/source/enumerator/DynamoDbStreamsSourceEnumerator.java
b/flink-connector-aws/flink-connector-dynamodb/src/main/java/org/apache/flink/connector/dynamodb/source/enumerator/DynamoDbStreamsSourceEnumerator.java
index c178611..0fb5390 100644
---
a/flink-connector-aws/flink-connector-dynamodb/src/main/java/org/apache/flink/connector/dynamodb/source/enumerator/DynamoDbStreamsSourceEnumerator.java
+++
b/flink-connector-aws/flink-connector-dynamodb/src/main/java/org/apache/flink/connector/dynamodb/source/enumerator/DynamoDbStreamsSourceEnumerator.java
@@ -28,7 +28,6 @@ import org.apache.flink.configuration.Configuration;
import
org.apache.flink.connector.dynamodb.source.config.DynamodbStreamsSourceConfigConstants.InitialPosition;
import
org.apache.flink.connector.dynamodb.source.enumerator.event.SplitsFinishedEvent;
import
org.apache.flink.connector.dynamodb.source.enumerator.event.SplitsFinishedEventContext;
-import
org.apache.flink.connector.dynamodb.source.enumerator.tracker.SplitGraphInconsistencyTracker;
import
org.apache.flink.connector.dynamodb.source.enumerator.tracker.SplitTracker;
import
org.apache.flink.connector.dynamodb.source.exception.DynamoDbStreamsSourceException;
import org.apache.flink.connector.dynamodb.source.proxy.StreamProxy;
@@ -38,7 +37,6 @@ import
org.apache.flink.connector.dynamodb.source.util.ListShardsResult;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import software.amazon.awssdk.services.dynamodb.model.Shard;
-import software.amazon.awssdk.services.dynamodb.model.StreamStatus;
import javax.annotation.Nullable;
@@ -54,7 +52,6 @@ import java.util.Map.Entry;
import java.util.Set;
import java.util.stream.Collectors;
-import static
org.apache.flink.connector.dynamodb.source.config.DynamodbStreamsSourceConfigConstants.DESCRIBE_STREAM_INCONSISTENCY_RESOLUTION_RETRY_COUNT;
import static
org.apache.flink.connector.dynamodb.source.config.DynamodbStreamsSourceConfigConstants.SHARD_DISCOVERY_INTERVAL;
import static
org.apache.flink.connector.dynamodb.source.config.DynamodbStreamsSourceConfigConstants.STREAM_INITIAL_POSITION;
@@ -181,10 +178,6 @@ public class DynamoDbStreamsSourceEnumerator
throw new DynamoDbStreamsSourceException("Failed to list shards.",
throwable);
}
- if (discoveredSplits.getInconsistencyDetected()) {
- return;
- }
-
splitTracker.addSplits(discoveredSplits.getShards());
splitTracker.cleanUpOldFinishedSplits(
discoveredSplits.getShards().stream()
@@ -200,42 +193,6 @@ public class DynamoDbStreamsSourceEnumerator
assignAllAvailableSplits();
}
- /**
- * This method tracks the discovered splits in a graph and if the graph
has inconsistencies, it
- * tries to resolve them using DescribeStream calls using the first
inconsistent node found in
- * the split graph.
- *
- * @param discoveredSplits splits discovered after calling DescribeStream
at the start of the
- * application or periodically.
- */
- private SplitGraphInconsistencyTracker
trackSplitsAndResolveInconsistencies(
- ListShardsResult discoveredSplits) {
- SplitGraphInconsistencyTracker splitGraphInconsistencyTracker =
- new SplitGraphInconsistencyTracker();
- splitGraphInconsistencyTracker.addNodes(discoveredSplits.getShards());
-
- // we don't want to do inconsistency checks for DISABLED streams
because there will be no
- // open child shard in DISABLED stream
- boolean streamDisabled =
discoveredSplits.getStreamStatus().equals(StreamStatus.DISABLED);
- int describeStreamInconsistencyResolutionCount =
-
sourceConfig.get(DESCRIBE_STREAM_INCONSISTENCY_RESOLUTION_RETRY_COUNT);
- for (int i = 0;
- i < describeStreamInconsistencyResolutionCount
- && !streamDisabled
- &&
splitGraphInconsistencyTracker.inconsistencyDetected();
- i++) {
- String earliestClosedLeafNodeId =
- splitGraphInconsistencyTracker.getEarliestClosedLeafNode();
- LOG.warn(
- "We have detected inconsistency with DescribeStream
output, resolving inconsistency with shardId: {}",
- earliestClosedLeafNodeId);
- ListShardsResult shardsToResolveInconsistencies =
- streamProxy.listShards(streamArn,
earliestClosedLeafNodeId);
-
splitGraphInconsistencyTracker.addNodes(shardsToResolveInconsistencies.getShards());
- }
- return splitGraphInconsistencyTracker;
- }
-
private void assignAllAvailableSplits() {
List<DynamoDbStreamsShardSplit> splitsAvailableForAssignment =
splitTracker.splitsAvailableForAssignment();
@@ -282,26 +239,7 @@ public class DynamoDbStreamsSourceEnumerator
* @return list of discovered splits
*/
private ListShardsResult discoverSplits() {
- ListShardsResult listShardsResult = streamProxy.listShards(streamArn,
null);
- SplitGraphInconsistencyTracker splitGraphInconsistencyTracker =
- trackSplitsAndResolveInconsistencies(listShardsResult);
-
- ListShardsResult discoveredSplits = new ListShardsResult();
- discoveredSplits.setStreamStatus(listShardsResult.getStreamStatus());
-
discoveredSplits.setInconsistencyDetected(listShardsResult.getInconsistencyDetected());
- List<Shard> shardList = new
ArrayList<>(splitGraphInconsistencyTracker.getNodes());
- // We do not throw an exception here and just return to let
SplitTracker process through the
- // splits it has not yet processed. This might be helpful for large
streams which see a lot
- // of
- // inconsistency issues.
- if (splitGraphInconsistencyTracker.inconsistencyDetected()) {
- LOG.error(
- "There are inconsistencies in DescribeStream which we were
not able to resolve. First leaf node on which inconsistency was detected:"
- +
splitGraphInconsistencyTracker.getEarliestClosedLeafNode());
- return discoveredSplits;
- }
- discoveredSplits.addShards(shardList);
- return discoveredSplits;
+ return streamProxy.listShards(streamArn, null);
}
private void assignSplitToSubtask(
diff --git
a/flink-connector-aws/flink-connector-dynamodb/src/main/java/org/apache/flink/connector/dynamodb/source/enumerator/tracker/SplitGraphInconsistencyTracker.java
b/flink-connector-aws/flink-connector-dynamodb/src/main/java/org/apache/flink/connector/dynamodb/source/enumerator/tracker/SplitGraphInconsistencyTracker.java
deleted file mode 100644
index 68306fe..0000000
---
a/flink-connector-aws/flink-connector-dynamodb/src/main/java/org/apache/flink/connector/dynamodb/source/enumerator/tracker/SplitGraphInconsistencyTracker.java
+++ /dev/null
@@ -1,116 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one
- * or more contributor license agreements. See the NOTICE file
- * distributed with this work for additional information
- * regarding copyright ownership. The ASF licenses this file
- * to you under the Apache License, Version 2.0 (the
- * "License"); you may not use this file except in compliance
- * with the License. You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-
-package org.apache.flink.connector.dynamodb.source.enumerator.tracker;
-
-import
org.apache.flink.connector.dynamodb.source.enumerator.DynamoDbStreamsSourceEnumerator;
-import org.apache.flink.connector.dynamodb.source.util.ShardUtils;
-
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-import software.amazon.awssdk.services.dynamodb.model.Shard;
-
-import java.util.Collection;
-import java.util.HashMap;
-import java.util.HashSet;
-import java.util.List;
-import java.util.Map;
-import java.util.Set;
-import java.util.TreeSet;
-
-/**
- * Class to track the state of shard graph created as part of describestream
operation. This will
- * track for inconsistent shards returned due to describestream operation.
Caller will have to call
- * describestream again with the chronologically first leaf node to resolve
inconsistencies.
- */
-public class SplitGraphInconsistencyTracker {
- private final TreeSet<String> closedLeafNodes;
- private final Map<String, Shard> nodes;
- private static final Logger LOG =
- LoggerFactory.getLogger(DynamoDbStreamsSourceEnumerator.class);
-
- public SplitGraphInconsistencyTracker() {
- nodes = new HashMap<>();
- closedLeafNodes = new TreeSet<>();
- }
-
- /**
- * Adds shards to the shard graph tracker. We first add shards and once
shards are added, we
- * remove the parent shards from closedLeafNodes. We need 2x for loops
here to make sure if the
- * shards are not ordered, parent and child can both reside in
closedLeafNodes
- *
- * @param shards Shards to add in the tracker
- */
- public void addNodes(List<Shard> shards) {
- for (Shard shard : shards) {
- addNode(shard);
- }
- for (Shard shard : shards) {
- removeParentFromClosedLeafNodes(shard);
- }
- }
-
- private void removeParentFromClosedLeafNodes(Shard shard) {
- if (shard.parentShardId() != null) {
- closedLeafNodes.remove(shard.parentShardId());
- }
- }
-
- private void addNode(Shard shard) {
- nodes.put(shard.shardId(), shard);
- if (shard.sequenceNumberRange().endingSequenceNumber() != null) {
- closedLeafNodes.add(shard.shardId());
- }
- }
-
- /**
- * If any shard id is more than 24 hours old and it does not have a child,
we can assume that
- * the child has been trimmed. This case is pretty common for large
streams and it doesn't make
- * sense to halt the application for another expensive describestream
operation.
- */
- public boolean inconsistencyDetected() {
- Set<String> closedLeafNodesCopy = new HashSet<>(closedLeafNodes);
- for (String closedLeafNodeId : closedLeafNodesCopy) {
- if
(ShardUtils.isShardOlderThanInconsistencyDetectionRetentionPeriod(
- closedLeafNodeId)) {
- LOG.warn(
- "Shard id: {} has no child and has been created more
than 24 hours ago. Not tracking it and its ancestors",
- closedLeafNodeId);
- removeExpiredInconsistentLeaves(closedLeafNodeId);
- }
- }
- return !closedLeafNodes.isEmpty();
- }
-
- public String getEarliestClosedLeafNode() {
- return closedLeafNodes.first();
- }
-
- public Collection<Shard> getNodes() {
- return nodes.values();
- }
-
- private void removeExpiredInconsistentLeaves(String shardId) {
- String currentShardId = shardId;
- while (currentShardId != null && nodes.containsKey(currentShardId)) {
- Shard currentShard = nodes.remove(currentShardId);
- closedLeafNodes.remove(currentShardId);
- currentShardId = currentShard.parentShardId();
- }
- }
-}
diff --git
a/flink-connector-aws/flink-connector-dynamodb/src/main/java/org/apache/flink/connector/dynamodb/source/util/ListShardsResult.java
b/flink-connector-aws/flink-connector-dynamodb/src/main/java/org/apache/flink/connector/dynamodb/source/util/ListShardsResult.java
index c50a2eb..e714c54 100644
---
a/flink-connector-aws/flink-connector-dynamodb/src/main/java/org/apache/flink/connector/dynamodb/source/util/ListShardsResult.java
+++
b/flink-connector-aws/flink-connector-dynamodb/src/main/java/org/apache/flink/connector/dynamodb/source/util/ListShardsResult.java
@@ -36,12 +36,10 @@ import java.util.Objects;
public class ListShardsResult {
private final List<Shard> shards;
private StreamStatus streamStatus;
- private boolean inconsistencyDetected;
public ListShardsResult() {
this.shards = new ArrayList<>();
this.streamStatus = StreamStatus.ENABLED;
- this.inconsistencyDetected = false;
}
public void addShards(List<Shard> shardList) {
@@ -52,10 +50,6 @@ public class ListShardsResult {
this.streamStatus = streamStatus;
}
- public void setInconsistencyDetected(boolean inconsistencyDetected) {
- this.inconsistencyDetected = inconsistencyDetected;
- }
-
public List<Shard> getShards() {
return this.shards;
}
@@ -64,10 +58,6 @@ public class ListShardsResult {
return this.streamStatus;
}
- public boolean getInconsistencyDetected() {
- return this.inconsistencyDetected;
- }
-
@Override
public boolean equals(Object o) {
if (this == o) {
@@ -78,13 +68,12 @@ public class ListShardsResult {
}
ListShardsResult that = (ListShardsResult) o;
return Objects.equals(shards, that.shards)
- && Objects.equals(streamStatus, that.getStreamStatus())
- && Objects.equals(inconsistencyDetected,
that.inconsistencyDetected);
+ && Objects.equals(streamStatus, that.getStreamStatus());
}
@Override
public int hashCode() {
- return Objects.hash(shards, streamStatus, inconsistencyDetected);
+ return Objects.hash(shards, streamStatus);
}
@Override
@@ -94,8 +83,6 @@ public class ListShardsResult {
+ shards
+ ", streamStatus="
+ streamStatus.toString()
- + ", inconsistencyDetected="
- + inconsistencyDetected
+ "}";
}
}
diff --git
a/flink-connector-aws/flink-connector-dynamodb/src/main/java/org/apache/flink/connector/dynamodb/source/util/ShardUtils.java
b/flink-connector-aws/flink-connector-dynamodb/src/main/java/org/apache/flink/connector/dynamodb/source/util/ShardUtils.java
index 88a56b6..0f6f123 100644
---
a/flink-connector-aws/flink-connector-dynamodb/src/main/java/org/apache/flink/connector/dynamodb/source/util/ShardUtils.java
+++
b/flink-connector-aws/flink-connector-dynamodb/src/main/java/org/apache/flink/connector/dynamodb/source/util/ShardUtils.java
@@ -19,7 +19,6 @@
package org.apache.flink.connector.dynamodb.source.util;
import org.apache.flink.annotation.Internal;
-import
org.apache.flink.connector.dynamodb.source.enumerator.tracker.SplitGraphInconsistencyTracker;
import
org.apache.flink.connector.dynamodb.source.enumerator.tracker.SplitTracker;
import java.time.Duration;
@@ -36,14 +35,6 @@ public class ShardUtils {
*/
private static final Duration DDB_STREAMS_MAX_RETENTION_PERIOD =
Duration.ofHours(48);
- /**
- * Maximum retention period for shards to be stored in {@link
SplitGraphInconsistencyTracker}.
- * DDB Streams records are expired after 24 hours. We do not need to
resolve inconsistencies
- * during shard expiry time.
- */
- private static final Duration
DDB_STREAMS_MAX_RETENTION_PERIOD_FOR_RESOLVING_INCONSISTENCIES =
- Duration.ofHours(25);
-
private static final String SHARD_ID_SEPARATOR = "-";
/**
@@ -66,20 +57,6 @@ public class ShardUtils {
.isAfter(getShardCreationTime(shardId).plus(DDB_STREAMS_MAX_RETENTION_PERIOD));
}
- /**
- * Returns true if the shard is older than what we want to store in {@link
- * SplitGraphInconsistencyTracker}.
- *
- * @param shardId
- */
- public static boolean
isShardOlderThanInconsistencyDetectionRetentionPeriod(String shardId) {
- return Instant.now()
- .isAfter(
- getShardCreationTime(shardId)
- .plus(
-
DDB_STREAMS_MAX_RETENTION_PERIOD_FOR_RESOLVING_INCONSISTENCIES));
- }
-
/**
* Returns true if the shard was created before the given timestamp.
*
diff --git
a/flink-connector-aws/flink-connector-dynamodb/src/test/java/org/apache/flink/connector/dynamodb/source/enumerator/tracker/SplitGraphInconsistencyTrackerTest.java
b/flink-connector-aws/flink-connector-dynamodb/src/test/java/org/apache/flink/connector/dynamodb/source/enumerator/tracker/SplitGraphInconsistencyTrackerTest.java
deleted file mode 100644
index 8127485..0000000
---
a/flink-connector-aws/flink-connector-dynamodb/src/test/java/org/apache/flink/connector/dynamodb/source/enumerator/tracker/SplitGraphInconsistencyTrackerTest.java
+++ /dev/null
@@ -1,96 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one
- * or more contributor license agreements. See the NOTICE file
- * distributed with this work for additional information
- * regarding copyright ownership. The ASF licenses this file
- * to you under the Apache License, Version 2.0 (the
- * "License"); you may not use this file except in compliance
- * with the License. You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-
-package org.apache.flink.connector.dynamodb.source.enumerator.tracker;
-
-import org.junit.jupiter.api.Test;
-import software.amazon.awssdk.services.dynamodb.model.Shard;
-
-import java.time.Duration;
-import java.time.Instant;
-import java.util.Arrays;
-import java.util.List;
-
-import static
org.apache.flink.connector.dynamodb.source.util.TestUtil.generateShard;
-import static
org.apache.flink.connector.dynamodb.source.util.TestUtil.generateShardId;
-import static org.assertj.core.api.Assertions.assertThat;
-
-/**
- * Tests the {@link SplitGraphInconsistencyTracker} class to verify that it
correctly discovers
- * inconsistencies within the shard graph.
- */
-public class SplitGraphInconsistencyTrackerTest {
- @Test
- public void testShardGraphTrackerHappyCase() {
- List<Shard> shards =
- Arrays.asList(
- // shards which don't have a parent
- generateShard(0, "1400", "1700", null),
- generateShard(1, "1500", "1800", null),
- // shards produced by rotation of parents
- generateShard(2, "1710", null, generateShardId(0)),
- generateShard(3, "1520", null, generateShardId(1)));
- SplitGraphInconsistencyTracker splitGraphInconsistencyTracker =
- new SplitGraphInconsistencyTracker();
- splitGraphInconsistencyTracker.addNodes(shards);
-
assertThat(splitGraphInconsistencyTracker.inconsistencyDetected()).isFalse();
- assertThat(splitGraphInconsistencyTracker.getNodes())
- .containsExactlyInAnyOrderElementsOf(shards);
- }
-
- @Test
- public void
-
testSplitGraphInconsistencyTrackerDetectsInconsistenciesAndDoesNotDeleteClosedLeafNode()
{
- long firstShardIdTimestamp = Instant.now().toEpochMilli();
- long secondShardIdTimestamp =
Instant.now().minus(Duration.ofHours(2)).toEpochMilli();
- List<Shard> shards =
- Arrays.asList(
- // shards which don't have a parent
- generateShard(firstShardIdTimestamp, "1400", "1700",
null),
- generateShard(secondShardIdTimestamp, "1500", "1800",
null),
- // shards produced by rotation of parents
- generateShard(2, "1710", null,
generateShardId(firstShardIdTimestamp)));
- SplitGraphInconsistencyTracker splitGraphInconsistencyTracker =
- new SplitGraphInconsistencyTracker();
- splitGraphInconsistencyTracker.addNodes(shards);
-
assertThat(splitGraphInconsistencyTracker.inconsistencyDetected()).isTrue();
- assertThat(splitGraphInconsistencyTracker.getEarliestClosedLeafNode())
- .isEqualTo(generateShardId(secondShardIdTimestamp));
- assertThat(splitGraphInconsistencyTracker.getNodes())
- .containsExactlyInAnyOrderElementsOf(shards);
- }
-
- @Test
- public void
testSplitGraphInconsistencyTrackerDetectsInconsistenciesAndDeletesClosedLeafNode()
{
- long firstShardIdTimestamp = Instant.now().toEpochMilli();
- long secondShardIdTimestamp =
Instant.now().minus(Duration.ofHours(26)).toEpochMilli();
- List<Shard> shards =
- Arrays.asList(
- // shards which don't have a parent
- generateShard(firstShardIdTimestamp, "1400", "1700",
null),
- generateShard(secondShardIdTimestamp, "1500", "1800",
null),
- // shards produced by rotation of parents
- generateShard(2, "1710", null,
generateShardId(firstShardIdTimestamp)));
- SplitGraphInconsistencyTracker splitGraphInconsistencyTracker =
- new SplitGraphInconsistencyTracker();
- splitGraphInconsistencyTracker.addNodes(shards);
-
assertThat(splitGraphInconsistencyTracker.inconsistencyDetected()).isFalse();
- assertThat(splitGraphInconsistencyTracker.getNodes())
- .containsExactlyInAnyOrder(shards.get(0), shards.get(2));
- }
-}
diff --git
a/flink-connector-aws/flink-connector-dynamodb/src/test/java/org/apache/flink/connector/dynamodb/source/util/ListShardsResultTest.java
b/flink-connector-aws/flink-connector-dynamodb/src/test/java/org/apache/flink/connector/dynamodb/source/util/ListShardsResultTest.java
index 1dc71c6..b3f5bb6 100644
---
a/flink-connector-aws/flink-connector-dynamodb/src/test/java/org/apache/flink/connector/dynamodb/source/util/ListShardsResultTest.java
+++
b/flink-connector-aws/flink-connector-dynamodb/src/test/java/org/apache/flink/connector/dynamodb/source/util/ListShardsResultTest.java
@@ -48,13 +48,6 @@ public class ListShardsResultTest {
assertThat(listShardsResult.getStreamStatus()).isEqualTo(StreamStatus.ENABLED);
}
- @Test
- void testSetInconsistencyDetected() {
- ListShardsResult listShardsResult = new ListShardsResult();
- listShardsResult.setInconsistencyDetected(true);
-
assertThat(listShardsResult.getInconsistencyDetected()).isEqualTo(true);
- }
-
@Test
void testEquals() {
EqualsVerifier.simple().forClass(ListShardsResult.class).verify();
diff --git
a/flink-connector-aws/flink-connector-dynamodb/src/test/java/org/apache/flink/connector/dynamodb/source/util/ShardUtilsTest.java
b/flink-connector-aws/flink-connector-dynamodb/src/test/java/org/apache/flink/connector/dynamodb/source/util/ShardUtilsTest.java
index 2d6de15..52cd1d3 100644
---
a/flink-connector-aws/flink-connector-dynamodb/src/test/java/org/apache/flink/connector/dynamodb/source/util/ShardUtilsTest.java
+++
b/flink-connector-aws/flink-connector-dynamodb/src/test/java/org/apache/flink/connector/dynamodb/source/util/ShardUtilsTest.java
@@ -42,13 +42,4 @@ public class ShardUtilsTest {
assertThat(ShardUtils.isShardOlderThanRetentionPeriod(oldShardId)).isTrue();
assertThat(ShardUtils.isShardOlderThanRetentionPeriod(newShardId)).isFalse();
}
-
- @Test
- void testIsShardOlderThanInconsistencyDetectionRetentionPeriod() {
- Instant currentTime = Instant.now();
- String oldShardId = "shardId-" +
currentTime.minus(OLD_SHARD_DURATION).toEpochMilli();
- String newShardId = "shardId-" + currentTime.toEpochMilli();
-
assertThat(ShardUtils.isShardOlderThanRetentionPeriod(oldShardId)).isTrue();
-
assertThat(ShardUtils.isShardOlderThanRetentionPeriod(newShardId)).isFalse();
- }
}