JingsongLi commented on code in PR #9535:
URL: https://github.com/apache/paimon/pull/9535#discussion_r4182572237


##########
paimon-core/src/main/java/org/apache/paimon/table/source/DataTableStreamScan.java:
##########
@@ -315,6 +317,35 @@ public Long watermark() {
 
     @Override
     public void restore(@Nullable Long nextSnapshotId) {
+        if (nextSnapshotId != null) {
+            Long earliestSnapshotId = snapshotManager.earliestSnapshotId();
+            if (earliestSnapshotId != null && earliestSnapshotId > 
nextSnapshotId) {
+                // The restored snapshot has already been expired. Whether the 
consumer can still
+                // resume from it depends on the changelog lifecycle.
+                if (options.changelogLifecycleDecoupled()
+                        && 
changelogManager.longLivedChangelogExists(nextSnapshotId)) {
+                    // The long-lived changelog outlives the snapshot, so keep 
reading from the
+                    // changelog instead of replaying data the consumer has 
already committed.
+                    LOG.warn(
+                            "The restored snapshot with id {} has expired, but 
its long-lived "
+                                    + "changelog is still available. Resuming 
from the changelog.",
+                            nextSnapshotId);
+                    this.nextSnapshotId = nextSnapshotId;
+                    return;
+                }
+                // No changelog to fall back on. Keeping the expired id would 
make every restart
+                // fail with OutOfRangeException, so the job could never 
self-recover; restart the
+                // scan from the starting scanner instead.
+                LOG.warn(
+                        "The restored snapshot with id {} has expired. "
+                                + "The earliest snapshot is {}. "
+                                + "Falling back to starting scanner.",
+                        nextSnapshotId,
+                        earliestSnapshotId);
+                this.nextSnapshotId = null;

Review Comment:
   [P1] Limit startup fallback to the dedicated compaction recovery contract
   
   The added long-lived-changelog branch preserves that recovery path, but 
ordinary consumers without long-lived changelog still hit this global reset. On 
a real latest-mode consumer, checkpoint() returned 2 after snapshot 1; 
subsequent commits expired 1/2 and retained 3/4. Restoring checkpoint 2 now 
clears it and reapplies ContinuousLatestStartingScanner, advances checkpoint to 
5, and emits zero rows from retained snapshots 3/4. An earliest-retained 
control reads both delta rows, while the original implementation retains the 
stalled checkpoint 2. Thus this converts an explicit recovery gap into silent 
skipping of still-available data. Scope the reset to 
ContinuousCompactorStartingScanner/the dedicated source, or introduce an 
explicit general recovery policy which does not reapply the initial startup 
mode. Add a restore regression with expired checkpoint plus later retained 
snapshots; the current existing startup tests all pass without exercising this 
case.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to