diff --git a/paimon-core/src/main/java/org/apache/paimon/utils/SnapshotManager.java b/paimon-core/src/main/java/org/apache/paimon/utils/SnapshotManager.java index 8cd159997865..45c6e8224ef8 100644 --- a/paimon-core/src/main/java/org/apache/paimon/utils/SnapshotManager.java +++ b/paimon-core/src/main/java/org/apache/paimon/utils/SnapshotManager.java @@ -382,8 +382,8 @@ public long repairEarliestSnapshot(long snapshotId) { earliest = mid + 1; // Search in the right half finalSnapshot = snapshot; } else { - finalSnapshot = snapshot; // Found the exact match - break; + finalSnapshot = snapshot; + earliest = mid + 1; } } return finalSnapshot; @@ -415,8 +415,8 @@ public long repairEarliestSnapshot(long snapshotId) { } else if (commitTime < timestampMills) { earliest = mid + 1; // Search in the right half } else { - finalSnapshot = snapshot; // Found the exact match - break; + finalSnapshot = snapshot; + latest = mid - 1; } } return finalSnapshot; diff --git a/paimon-core/src/test/java/org/apache/paimon/utils/SnapshotManagerTest.java b/paimon-core/src/test/java/org/apache/paimon/utils/SnapshotManagerTest.java index bf1a005878a9..fbdb033a7fab 100644 --- a/paimon-core/src/test/java/org/apache/paimon/utils/SnapshotManagerTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/utils/SnapshotManagerTest.java @@ -423,6 +423,34 @@ public void testLaterOrEqualTimeMills() throws IOException { assertThat(snapshotManager.laterOrEqualTimeMills(millis + 10001)).isNull(); } + @Test + public void testEarlierOrEqualTimeMillsWithDuplicateCommitTimes() throws IOException { + long millis = 1684726826L; + FileIO localFileIO = LocalFileIO.create(); + SnapshotManager snapshotManager = + newSnapshotManager(localFileIO, new Path(tempDir.toString())); + for (long i = 0; i < 3; i++) { + Snapshot snapshot = createSnapshotWithMillis(i, millis); + localFileIO.tryToWriteAtomic(snapshotManager.snapshotPath(i), snapshot.toJson()); + } + + assertThat(snapshotManager.earlierOrEqualTimeMills(millis).id()).isEqualTo(2); + } + + @Test + public void testLaterOrEqualTimeMillsWithDuplicateCommitTimes() throws IOException { + long millis = 1684726826L; + FileIO localFileIO = LocalFileIO.create(); + SnapshotManager snapshotManager = + newSnapshotManager(localFileIO, new Path(tempDir.toString())); + for (long i = 0; i < 3; i++) { + Snapshot snapshot = createSnapshotWithMillis(i, millis); + localFileIO.tryToWriteAtomic(snapshotManager.snapshotPath(i), snapshot.toJson()); + } + + assertThat(snapshotManager.laterOrEqualTimeMills(millis).id()).isEqualTo(0); + } + @ParameterizedTest @ValueSource(booleans = {true, false}) public void testLaterOrEqualWatermark(boolean isRaceCondition) throws IOException {