This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git
The following commit(s) were added to refs/heads/master by this push:
new 45895a719e [core][python] Improve earliest snapshot retry handling
(#8242)
45895a719e is described below
commit 45895a719e6e3a92e1455f5eb9ccc3eb6096b967
Author: Jingsong Lee <[email protected]>
AuthorDate: Tue Jun 16 08:42:48 2026 +0800
[core][python] Improve earliest snapshot retry handling (#8242)
When the earliest snapshot is expired concurrently, callers may observe
an earliest snapshot id whose snapshot file is already gone. This PR
makes Java and Python retry forward through a bounded window before
giving up, while preserving empty-table handling for callers that have
no available latest snapshot.
---
.../table/source/snapshot/TimeTravelUtil.java | 22 +---------
.../org/apache/paimon/utils/SnapshotManager.java | 24 +++++++----
.../apache/paimon/utils/SnapshotManagerTest.java | 47 +++++++++++++-------
.../pypaimon/snapshot/snapshot_manager.py | 35 ++++++++++-----
.../pypaimon/tests/snapshot_manager_test.py | 50 ++++++++++++++++++++++
5 files changed, 124 insertions(+), 54 deletions(-)
diff --git
a/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/TimeTravelUtil.java
b/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/TimeTravelUtil.java
index 8d4087a4e7..28eceb69ce 100644
---
a/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/TimeTravelUtil.java
+++
b/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/TimeTravelUtil.java
@@ -48,7 +48,7 @@ import static
org.apache.paimon.CoreOptions.SCAN_TIMESTAMP_MILLIS;
import static org.apache.paimon.CoreOptions.SCAN_WATERMARK;
import static org.apache.paimon.utils.DateTimeUtils.parseTimestampData;
import static org.apache.paimon.utils.Preconditions.checkArgument;
-import static
org.apache.paimon.utils.SnapshotManager.EARLIEST_SNAPSHOT_DEFAULT_RETRY_NUM;
+import static org.apache.paimon.utils.SnapshotManager.retryEarliestSnapshot;
/** The util class of resolve snapshot from scan params for time travel. */
public class TimeTravelUtil {
@@ -223,29 +223,11 @@ public class TimeTravelUtil {
return null;
}
- if (stopSnapshotId == null) {
- stopSnapshotId = snapshotId + EARLIEST_SNAPSHOT_DEFAULT_RETRY_NUM;
- }
-
FunctionWithException<Long, Snapshot, FileNotFoundException>
snapshotFunction =
includeChangelog
? s -> tryGetChangelogOrSnapshot(snapshotManager,
changelogManager, s)
: snapshotManager::tryGetSnapshot;
-
- do {
- try {
- return snapshotFunction.apply(snapshotId);
- } catch (FileNotFoundException e) {
- snapshotId++;
- if (snapshotId > stopSnapshotId) {
- return null;
- }
- LOG.warn(
- "The earliest snapshot or changelog was once
identified but disappeared. "
- + "It might have been expired by other jobs
operating on this table. "
- + "Searching for the second earliest snapshot
or changelog instead. ");
- }
- } while (true);
+ return retryEarliestSnapshot(snapshotId, stopSnapshotId,
snapshotFunction);
}
private static Snapshot tryGetChangelogOrSnapshot(
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 fcc936460d..af21181214 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
@@ -65,7 +65,7 @@ public class SnapshotManager implements Serializable {
public static final String SNAPSHOT_PREFIX = "snapshot-";
- public static final int EARLIEST_SNAPSHOT_DEFAULT_RETRY_NUM = 3;
+ public static final int EARLIEST_SNAPSHOT_DEFAULT_RETRY_NUM = 300;
private final FileIO fileIO;
private final Path tablePath;
@@ -232,22 +232,30 @@ public class SnapshotManager implements Serializable {
return null;
}
+ return retryEarliestSnapshot(snapshotId, stopSnapshotId,
this::tryGetSnapshot);
+ }
+
+ public static Snapshot retryEarliestSnapshot(
+ long earliestSnapshotId,
+ @Nullable Long stopSnapshotId,
+ FunctionWithException<Long, Snapshot, FileNotFoundException>
snapshotFunction) {
if (stopSnapshotId == null) {
- stopSnapshotId = snapshotId + EARLIEST_SNAPSHOT_DEFAULT_RETRY_NUM;
+ stopSnapshotId = earliestSnapshotId +
EARLIEST_SNAPSHOT_DEFAULT_RETRY_NUM;
}
+ long snapshotId = earliestSnapshotId;
do {
try {
- return tryGetSnapshot(snapshotId);
+ return snapshotFunction.apply(snapshotId);
} catch (FileNotFoundException e) {
snapshotId++;
if (snapshotId > stopSnapshotId) {
- return null;
+ throw new RuntimeException(
+ String.format(
+ "Cannot find earliest snapshot from #%s to
#%s.",
+ earliestSnapshotId, stopSnapshotId),
+ e);
}
- LOG.warn(
- "The earliest snapshot or changelog was once
identified but disappeared. "
- + "It might have been expired by other jobs
operating on this table. "
- + "Searching for the second earliest snapshot
or changelog instead. ");
}
} while (true);
}
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 486f0b3fe5..8e47b74aa3 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
@@ -93,6 +93,19 @@ public class SnapshotManagerTest {
assertThat(snapshotManager.earliestSnapshot().id()).isEqualTo(isRaceCondition ?
1 : 0);
}
+ @Test
+ public void testEarliestSnapshotThrowsWhenRetryExhausted() throws
IOException {
+ FileIO localFileIO = LocalFileIO.create();
+ SnapshotManager snapshotManager =
+ new TestSnapshotManager(localFileIO, new
Path(tempDir.toString()), true);
+ Snapshot snapshot = createSnapshotWithMillis(0, 1684726826L);
+ localFileIO.tryToWriteAtomic(snapshotManager.snapshotPath(0),
snapshot.toJson());
+
+ assertThatThrownBy(snapshotManager::earliestSnapshot)
+ .isInstanceOf(RuntimeException.class)
+ .hasMessageContaining("Cannot find earliest snapshot");
+ }
+
@ParameterizedTest
@ValueSource(booleans = {true, false})
public void testEarlierOrEqualWatermark(boolean isRaceCondition) throws
IOException {
@@ -116,7 +129,7 @@ public class SnapshotManagerTest {
long base = System.currentTimeMillis();
ThreadLocalRandom random = ThreadLocalRandom.current();
- int numSnapshots = random.nextInt(1, 20);
+ int numSnapshots = isRaceCondition ? random.nextInt(2, 20) :
random.nextInt(1, 20);
Set<Long> set = new HashSet<>();
while (set.size() < numSnapshots) {
set.add(base + random.nextLong(0, 1_000_000));
@@ -146,25 +159,13 @@ public class SnapshotManagerTest {
TimeTravelUtil.earlierThanTimeMills(snapshotManager, null,
time, false, false);
if (millis.get(numSnapshots - 1) < time) {
- if (isRaceCondition && millis.size() == 1) {
- if (tries == 0) {
-
assertThat(actual).isLessThanOrEqualTo(firstSnapshotId);
- } else {
- assertThat(actual).isNull();
- }
- } else {
- assertThat(actual).isEqualTo(firstSnapshotId +
numSnapshots - 1);
- }
+ assertThat(actual).isEqualTo(firstSnapshotId + numSnapshots -
1);
} else {
for (int i = 0; i < numSnapshots; i++) {
if (millis.get(i) >= time) {
if (isRaceCondition && i == 0) {
// The first snapshot expired during invocation
- if (millis.size() == 1 && tries > 0) {
- assertThat(actual).isNull();
- } else {
-
assertThat(actual).isLessThanOrEqualTo(firstSnapshotId);
- }
+
assertThat(actual).isLessThanOrEqualTo(firstSnapshotId);
} else {
assertThat(actual).isLessThanOrEqualTo(firstSnapshotId + i - 1);
}
@@ -175,6 +176,22 @@ public class SnapshotManagerTest {
}
}
+ @Test
+ public void testEarlierThanTimeMillisThrowsWhenEarliestRetryExhausted()
throws IOException {
+ FileIO localFileIO = LocalFileIO.create();
+ SnapshotManager snapshotManager =
+ new TestSnapshotManager(localFileIO, new
Path(tempDir.toString()), true);
+ Snapshot snapshot = createSnapshotWithMillis(0, 1684726826L);
+ localFileIO.tryToWriteAtomic(snapshotManager.snapshotPath(0),
snapshot.toJson());
+
+ assertThatThrownBy(
+ () ->
+ TimeTravelUtil.earlierThanTimeMills(
+ snapshotManager, null, 1684726827L,
false, false))
+ .isInstanceOf(RuntimeException.class)
+ .hasMessageContaining("Cannot find earliest snapshot");
+ }
+
@ParameterizedTest
@ValueSource(booleans = {true, false})
public void testEarlierOrEqualTimeMills(boolean isRaceCondition) throws
IOException {
diff --git a/paimon-python/pypaimon/snapshot/snapshot_manager.py
b/paimon-python/pypaimon/snapshot/snapshot_manager.py
index 41a1ac81ba..d404a0dd94 100644
--- a/paimon-python/pypaimon/snapshot/snapshot_manager.py
+++ b/paimon-python/pypaimon/snapshot/snapshot_manager.py
@@ -26,6 +26,8 @@ from pypaimon.common.json_util import JSON
from pypaimon.snapshot.snapshot import Snapshot
from pypaimon.snapshot.snapshot_loader import SnapshotLoader
+EARLIEST_SNAPSHOT_DEFAULT_RETRY_NUM = 300
+
class SnapshotManager:
"""Manager for snapshot files using unified FileIO."""
@@ -158,19 +160,27 @@ class SnapshotManager:
"""Get the file path for the given snapshot ID."""
return f"{self.snapshot_dir}/snapshot-{snapshot_id}"
- def try_get_earliest_snapshot(self) -> Optional[Snapshot]:
+ def try_get_earliest_snapshot(
+ self, stop_snapshot_id: Optional[int] = None
+ ) -> Optional[Snapshot]:
earliest_file = f"{self.snapshot_dir}/EARLIEST"
if self.file_io.exists(earliest_file):
earliest_content = self.file_io.read_file_utf8(earliest_file)
earliest_snapshot_id = int(earliest_content.strip())
- snapshot = self.get_snapshot_by_id(earliest_snapshot_id)
- if snapshot is None:
- logger.warning(
- "The earliest snapshot or changelog was once identified
but disappeared. "
- "It might have been expired by other jobs operating on
this table."
- )
- return snapshot
- return self.get_snapshot_by_id(1)
+ else:
+ earliest_snapshot_id = 1
+
+ if stop_snapshot_id is None:
+ stop_snapshot_id = earliest_snapshot_id +
EARLIEST_SNAPSHOT_DEFAULT_RETRY_NUM
+
+ for snapshot_id in range(earliest_snapshot_id, stop_snapshot_id + 1):
+ snapshot = self.get_snapshot_by_id(snapshot_id)
+ if snapshot is not None:
+ return snapshot
+ raise RuntimeError(
+ f"Cannot find earliest snapshot from #{earliest_snapshot_id} "
+ f"to #{stop_snapshot_id}."
+ )
def earlier_or_equal_time_mills(self, timestamp: int) ->
Optional[Snapshot]:
"""
@@ -182,10 +192,13 @@ class SnapshotManager:
Returns:
The latest snapshot with time_millis <= timestamp, or None if no
such snapshot exists
"""
- earliest_snap = self.try_get_earliest_snapshot()
latest_snap = self.get_latest_snapshot()
- if earliest_snap is None or latest_snap is None:
+ if latest_snap is None:
+ return None
+
+ earliest_snap = self.try_get_earliest_snapshot(latest_snap.id)
+ if earliest_snap is None:
return None
earliest = earliest_snap.id
diff --git a/paimon-python/pypaimon/tests/snapshot_manager_test.py
b/paimon-python/pypaimon/tests/snapshot_manager_test.py
index 31dbfa84e6..c5adcb0537 100644
--- a/paimon-python/pypaimon/tests/snapshot_manager_test.py
+++ b/paimon-python/pypaimon/tests/snapshot_manager_test.py
@@ -33,6 +33,12 @@ def _create_mock_snapshot(snapshot_id: int, commit_kind: str
= "APPEND"):
return snapshot
+def _create_mock_snapshot_with_time(snapshot_id: int, time_millis: int):
+ snapshot = _create_mock_snapshot(snapshot_id)
+ snapshot.time_millis = time_millis
+ return snapshot
+
+
def _build_manager(file_io):
from pypaimon.snapshot.snapshot_manager import SnapshotManager
return SnapshotManager(file_io, "/tmp/test_table")
@@ -111,6 +117,50 @@ class SnapshotManagerTest(unittest.TestCase):
self.assertEqual(next_id, 8)
self.assertEqual(skipped_count, 3)
+ def test_earlier_or_equal_time_mills_skips_missing_earliest_snapshot(self):
+ """earlier_or_equal_time_mills should retry if earliest snapshot
disappeared."""
+ file_io = Mock()
+ file_io.exists.return_value = True
+ file_io.read_file_utf8.return_value = "1"
+
+ snapshots = {
+ 2: _create_mock_snapshot_with_time(2, 2000),
+ 3: _create_mock_snapshot_with_time(3, 3000),
+ }
+
+ manager = _build_manager(file_io)
+ manager.get_latest_snapshot = lambda: snapshots[3]
+ manager.get_snapshot_by_id = lambda sid: snapshots.get(sid)
+
+ result = manager.earlier_or_equal_time_mills(2500)
+
+ self.assertEqual(result.id, 2)
+
+ def
test_try_get_earliest_snapshot_retries_beyond_one_hundred_missing_snapshots(self):
+ file_io = Mock()
+ file_io.exists.return_value = True
+ file_io.read_file_utf8.return_value = "1"
+
+ snapshot = _create_mock_snapshot_with_time(105, 105000)
+
+ manager = _build_manager(file_io)
+ manager.get_snapshot_by_id = lambda sid: snapshot if sid == 105 else
None
+
+ result = manager.try_get_earliest_snapshot()
+
+ self.assertEqual(result.id, 105)
+
+ def test_try_get_earliest_snapshot_throws_when_retry_exhausted(self):
+ file_io = Mock()
+ file_io.exists.return_value = True
+ file_io.read_file_utf8.return_value = "1"
+
+ manager = _build_manager(file_io)
+ manager.get_snapshot_by_id = lambda sid: None
+
+ with self.assertRaisesRegex(RuntimeError, "Cannot find earliest
snapshot"):
+ manager.try_get_earliest_snapshot()
+
if __name__ == '__main__':
unittest.main()