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]