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

Reply via email to