zhang-arvin commented on PR #9535: URL: https://github.com/apache/paimon/pull/9535#issuecomment-5694541024
Thanks @lilei1128 for the thorough question. I traced the restore path to verify: **Pending splits:** `DataTableStreamScan.restore(nextSnapshotId)` is the single restore entry point for all scan modes — the Flink source restores it first, then re-enters `nextPlan()`. So the fallback in `restore()` covers the pending-splits case: after fallback the scanner restarts from `earliestSnapshotId`, and splits re-generated from that point carry valid snapshot ids. Split *consumption* does not re-look up `snapshotManager.snapshot(snapshotId)` metadata — `DataTableRead.executePlan` reads the `DataFileMeta` paths carried in the split directly (see `DataSplit` holding `dataFiles`/`beforeFiles` and `KeyValueTableRead` reading via file paths). **Sink/writer state:** that lives in the Flink sink path, which this source-side fix does not touch. If the dedicated compact job's writer checkpoint can reference an expired snapshot id, that would be a separate recovery gap worth tracking in its own issue rather than expanding this PR. I'll rebase onto latest master so the Spark test failure (build_test on the older head) is re-verified on the current base. -- 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]
