From f2386a0f62e5c9280e05a8943906700e0ddba8b9 Mon Sep 17 00:00:00 2001 From: zhang-arvin Date: Tue, 1 Sep 2026 23:23:04 +0800 Subject: [PATCH] [Bug] Fix Flink Dedicated Streaming Compact Job stuck when snapshot expired --- .../paimon/table/source/DataTableStreamScan.java | 13 +++++++++++++ 1 file changed, 13 insertions(+) diff --git a/paimon-core/src/main/java/org/apache/paimon/table/source/DataTableStreamScan.java b/paimon-core/src/main/java/org/apache/paimon/table/source/DataTableStreamScan.java index 02806cfe33fa..ccdd3039d005 100644 --- a/paimon-core/src/main/java/org/apache/paimon/table/source/DataTableStreamScan.java +++ b/paimon-core/src/main/java/org/apache/paimon/table/source/DataTableStreamScan.java @@ -315,6 +315,19 @@ public Long watermark() { @Override public void restore(@Nullable Long nextSnapshotId) { + if (nextSnapshotId != null) { + Long earliestSnapshotId = snapshotManager.earliestSnapshotId(); + if (earliestSnapshotId != null && earliestSnapshotId > nextSnapshotId) { + LOG.warn( + "The restored snapshot with id {} has expired. " + + "The earliest snapshot is {}. " + + "Falling back to starting scanner.", + nextSnapshotId, + earliestSnapshotId); + this.nextSnapshotId = null; + return; + } + } this.nextSnapshotId = nextSnapshotId; }