Skip to content

[core] Fall back to listing when the LATEST hint points to an expired snapshot - #10323

Open
dev-donghwan wants to merge 2 commits into
apache:masterfrom
dev-donghwan:fix-latest-hint-stale
Open

dev-donghwan wants to merge 2 commits into
apache:masterfrom
dev-donghwan:fix-latest-hint-stale

Conversation

@dev-donghwan

@dev-donghwan dev-donghwan commented Sep 30, 2026 •

Copy link
Copy Markdown
Contributor

Related: #10324 (merged) fixes one way the LATEST hint falls behind (a snapshot rename that throws but has succeeded).

Purpose

A stale LATEST hint should be harmless, because the hint is only a cache. Today it can leave a table permanently unable to commit or read, even though all newer snapshot files are still there.

The problem

HintFileUtils.findLatest reads the hint N and returns it if snapshot-(N + 1) does not exist. It never checks that snapshot-N itself exists:

Long snapshotId = readHint(fileIO, LATEST, dir);
if (snapshotId != null && snapshotId > 0) {
    long nextSnapshot = snapshotId + 1;
    // it is the latest only there is no next one
    if (!fileIO.exists(file.apply(nextSnapshot))) {
        return snapshotId;
    }
}
return findByListFiles(fileIO, Math::max, dir, prefix);

A missing N + 1 can mean two things: it has not been created yet, or it was created and then expired. findLatest assumes the first. If the hint stops moving while commits and expiration go on, expiration eventually removes both N and N + 1:

snapshot-13 ... snapshot-30    (10, 11 and 12 expired)
LATEST = 10

findLatest sees that snapshot-11 is missing and returns 10, which no longer exists, instead of 30. The directory listing is never reached.

The table cannot recover on its own:

  1. Only the commit path writes LATEST.
  2. A commit first calls latestSnapshot(), which gets 10 from findLatest and fails with Snapshot file .../snapshot-10 does not exist. It might have been expired by other jobs operating on this table. ....
  3. The commit never gets far enough to rewrite the hint, so every later commit reads the same stale hint.

The retry added in feabf2c does not help, because the second findLatest call returns the same id. Callers that use only latestSnapshotId(), such as expiration and the streaming starting scanners, get the wrong id as well.

Why this was safe before

findLatest and findEarliest were introduced together in FLINK-26778 (#56). Each hint only guards against the gap its own writer can leave:

EARLIEST LATEST
Written by expiration commit
Order delete snapshots, then write hint create snapshot, then write hint
Gap in between hint points to a deleted snapshot a newer snapshot exists
Check snapshot-N exists snapshot-(N + 1) does not exist

The review of #56 assumed that LATEST "is only changed by the commit operation and that value should be precise". A failed hint write also failed the commit. Under that assumption, expiration could never reach the snapshot LATEST points to, so checking that it exists was unnecessary.

That assumption no longer always holds. Since #5771 (1.3.0), a commit retry that finds its snapshot already written treats the commit as successful. This is the "Check if the commit has been completed" path in FileStoreCommitImpl. It correctly avoids a duplicate commit, but it does not rewrite the hint. So LATEST can now stay behind across many successful commits. Once it falls behind by more than the retention window, the hinted snapshot is expired.

Writing the hint on that path would help, but it cannot cover hint writes that keep failing. This PR therefore makes the table recover from a stale hint, whatever made it stale. #10324 separately fixes the case where a snapshot rename throws but has succeeded, which is how the hint fell behind for us.

How we hit this

We run Paimon 1.4.2 on Flink 2.2.1, with a Hive catalog and the warehouse on an S3-compatible object store. The table keeps snapshot.num-retained.max=20 with 30s checkpoints, which is about 6 minutes of snapshots.

The trigger was a problem in our own setup. S3A was loaded from /opt/flink/lib rather than as a Flink plugin. After a job restart closed the user classloader, S3A copies on a long-lived TaskManager kept succeeding on the server but failing on the client side. Every snapshot rename threw even though the snapshot file was written, and the retry path above reported success without writing the hint. LATEST stayed at the same id for 21 consecutive commits. We have since fixed this by moving S3A into the plugin directory.

That trigger alone should have been temporary. The permanent outage came from findLatest. About 6 minutes after the hint stopped moving, expiration removed the hinted snapshot. From then on, the writer failed in FileStoreCommitImpl.tryCommit and in FileSystemWriteRestore.restoreFiles, and a downstream streaming reader failed with OutOfRangeException. The job restarted about 6,000 times over two days. Overwriting LATEST by hand was the only way out.

Any failure that keeps LATEST from moving for longer than the retention window leads to the same state, for example repeated hint write failures, or a filesystem that keeps reporting errors after successful renames. With a small snapshot.num-retained.max, that window can be only a few minutes.

Change

findLatest is unchanged, so the normal path does not get any extra IO.

The fix is in the existing fallback of SnapshotManager#latestSnapshotFromFileSystem. When reading the hinted snapshot fails with FileNotFoundException and findLatest still returns the same id, it now lists the snapshot directory to find the real latest snapshot and reads that one. If the listing finds no newer snapshot, it throws as before.

State of hint N Before After
N is the latest read N read N (no extra IO)
N + 1 exists (hint slightly behind) findLatest lists same
N and N + 1 both expired read N fails, retry returns N, throws read N fails, retry returns N, list and read the real latest

The commit path (FileStoreCommitImpl#tryCommit) and the writer restore (FileSystemWriteRestore#restoreFiles) both go through this method. They now find the real latest snapshot, the next commit succeeds and rewrites the hint, and the table recovers by itself. Callers that only use latestSnapshotId() may still see the stale id until that commit rewrites the hint.

Tests

  • SnapshotManagerTest#testLatestSnapshotWithExpiredLatestHint: the hint points to an expired snapshot while snapshots 5–10 exist. latestSnapshot() threw before this change and returns snapshot 10 after it.
  • StaleLatestHintTest#testCommitAfterLatestHintExpired: drives the real commit and expire paths with a FileIO that skips the LATEST write, then lets hint writes recover. Before this change, commits and reads keep failing after the recovery. After it, the next commit succeeds and refreshes the hint.
  • Existing SnapshotManagerTest cases pass, including testLatestSnapshotStillFailsWhenNoNewerSnapshotExists.
  • All tests under org.apache.paimon.utils, org.apache.paimon.table.source and org.apache.paimon.catalog, plus the expiration and commit tests in org.apache.paimon.operation: 822 tests pass.
  • Spotless and Checkstyle for paimon-core

… snapshot

findLatest trusted the LATEST hint whenever snapshot-(hint + 1) did not
exist, without checking that the hinted snapshot itself still exists.
When the hint stays behind across more commits than the retention
window, expiration removes both the hinted snapshot and the next one,
and findLatest keeps returning the expired id. Commits then fail before
they can rewrite the hint, so the table never recovers by itself.

Also require the hinted snapshot to exist, as findEarliest already does,
and fall back to listing otherwise.

@JingsongLi JingsongLi left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

-1 here will another io.

…issing

Address review: do not add an exists check to findLatest, which would
cost another IO on every lookup. Keep findLatest as it is, and list the
snapshot directory only in latestSnapshotFromFileSystem after the
hinted snapshot turns out to be missing and the hint still returns the
same id. The commit path then finds the real latest snapshot and
rewrites the hint, so the table recovers without extra IO on the
normal path.
@dev-donghwan

Copy link
Copy Markdown
Contributor Author

Thanks for the review, @JingsongLi! You're right, I overlooked that findLatest is on the hot path and the extra exists check would cost another IO on every lookup.

I've reverted the findLatest change and moved the fix into the existing fallback in latestSnapshotFromFileSystem. Only when reading the hinted snapshot fails with FileNotFoundException and the hint still returns the same id, it now lists the snapshot directory to find the real latest snapshot. If the listing finds no newer snapshot, it still throws as before. So there's no extra IO on the normal path, and the commit path can still recover and rewrite the hint.

Could you take another look when you have time?

@JingsongLi

Copy link
Copy Markdown
Contributor

Reviewed head d8ef8b9aa0. The stale-hint outage has clear production value, but this revision introduces a replay correctness problem.

[P1] Recover commit-user deduplication before allowing the commit to proceed (SnapshotManager.java:227–232)

filterCommitted still calls latestSnapshotOfUser, which starts from latestSnapshotId(). If LATEST points to an expired snapshot and its successor is also expired, that lookup stops at the missing snapshot and reports no previous commit. The new fallback then lets tryCommit find the real latest snapshot and publish the same files again.

This is reachable in Flink's batch committer restart path: CommitterOperator#commitUpToCheckpoint(END_INPUT_CHECKPOINT_ID) calls filterAndCommit(committables, false, true), deliberately disabling append-file conflict checks after deduplication.

I reproduced this with an append table (bucket=-1), retained snapshot 7 containing the batch user's Long.MAX_VALUE commit, expired snapshots 1/2, and LATEST reset to 1. After reopening the table and replaying the original messages through filterAndCommitMultiple(..., false), this head returns 1, publishes snapshot 8, and a real table read returns 8 rows instead of 7, including the same row twice. With the baseline implementation, the same replay fails on the missing hinted snapshot before publishing. With append-file checks enabled, this head instead throws a duplicate-file conflict and cannot recover.

Please make the commit-user lookup recover from the missing hint as well before concluding that a commit is new, and add replay coverage for both conflict-check modes. A failure-only fallback can preserve the normal-path I/O cost.

There is also a remaining scope gap: latestSnapshotId() remains stale. Before another successful writer repairs LATEST, a fresh scan.mode=latest stream chose checkpoint 2 while the actual latest was 6 (expected checkpoint 7). The added test only checks reading after a successful commit repairs the hint. Please cover reader startup before that repair; this is an existing outage left unresolved, separate from the new duplicate-row regression.

Validation: SnapshotManagerTest + StaleLatestHintTest: 45 tests passed, including a final run without fast-build (Checkstyle/Spotless/enforcer enabled). Additional local replay/read probes reproduced the failures above; no probe changes were pushed.

@dev-donghwan

Copy link
Copy Markdown
Contributor Author

Thanks for the detailed review, @JingsongLi. You're right about both points. I reproduced the duplicate commit on d8ef8b9, both in core (replay with and without the append-file check) and through the Flink END_INPUT path on a bucket=-1 table.

I think I took the wrong approach with this PR, so before pushing anything I'd like to ask for your view as the original author of the hint files.

The current direction is to fix the read side. To continue with it, every caller that relies on the latest snapshot id would have to handle a hint that points to an expired snapshot. latestSnapshotId() alone has about 50 callers in the main code, and latestSnapshot() and latestSnapshotOfUser() have more, in core, Flink and Spark. I started with the ones you pointed out and fixed the dedup path, including the commit.last-safe-snapshot branch, which removed the duplicate commit. But the tests then showed other callers behaving differently from what they intend, because only latestSnapshot() gets the real latest while latestSnapshotId() still returns the stale id. For example:

  • incremental-between-timestamp silently reads all retained snapshots instead of the requested range (5 rows instead of 2), where master fails.
  • scan.mode=latest still starts from the stale id, so your second point stays unfixed.
  • The Flink/Spark rollback procedures call latestSnapshot() before rollbackTo, which is what currently stops them on a stale hint. With read-side recovery that call succeeds, and they go on to a rollbackTo that keeps the snapshots after the target and deletes newer tags.

Fixing each of these one by one would make the change much larger, and each fix could bring another side effect like these. So instead of the read side, I looked at the side that deletes snapshots. The stale state only appears when expiration deletes the snapshot the LATEST hint points to, so I'd like to propose preventing that instead.

Proposal

Snapshot expiration does not expire the snapshot the LATEST hint points to, similar to how it already keeps the snapshots consumers still need. Then findLatest works as it is. When the hint is behind, snapshot-(N + 1) exists, so it lists the directory and returns the real latest. All read-side changes are reverted, so every caller behaves exactly as on master.

  • No extra IO on the normal path. Expiration reads the LATEST hint once, and only when it is about to delete snapshots.
  • Expiration logs a warning whenever it keeps snapshots because the hint is behind.
  • I would skip the check when the catalog provides the latest snapshot itself (e.g. REST), where the LATEST file is not the source of truth.
  • A table that is already in the stale state behaves as on master. Some operations there are already wrong on master today (see below), so I think such tables are better handled by an explicit repair_latest_snapshot procedure, similar to repair_earliest_snapshot ([core][flink][spark] Support repairing the earliest snapshot hint #8883), as a follow-up.

If hint writes keep failing

In our incident, hint writes failed because of a broken TaskManager, which Paimon cannot prevent. What Paimon can avoid is turning that into a permanent outage: the TaskManager problem lasted about 6 minutes, while the table stayed stuck for two days until LATEST was fixed by hand.

I ran 30 commits with every LATEST write failing. With the proposal, all of them succeeded (each one slower because of the existing hint-write retries), reads and scan.mode=latest stayed correct, and snapshots piled up to 31. The first commit after hint writes recovered moved the hint, and expiration cleaned up back to the retention. On master the same run gets stuck once the hinted snapshot is expired.

Today the only trace of this is the generic Retry commit for exception warning, and the retry path that finds the commit already done logs nothing. If you agree, I'd also add a warning there, e.g. "Snapshot #N was committed by a previous attempt that failed, the LATEST hint may not have been updated", so the problem shows up at the first commit rather than only once snapshots pile up.

I also tried letting expiration rewrite the LATEST hint before deleting. It repairs the hint when another process runs expiration, but without a compare-and-swap on hint files it can race with rollback and write a hint that points to a snapshot the rollback has just deleted, which is the same stuck state. It also adds the hint-write retry delay to every expiration while writes fail. So I kept the proposal read-only.

Measured comparison

  • master: current behaviour
  • A: read-side recovery (d8ef8b9 plus the dedup fix)
  • B: the proposal above

All cells come from running the same probe on the four variants, except the rollback procedures, which are from reading the code.

1. A table whose LATEST hint stops moving while commits and expiration go on

master A B A + B
Hinted snapshot Expired Expired Kept Kept
New commit Fails on every attempt until LATEST is fixed by hand Succeeds Succeeds Succeeds
Committer replay (with/without append-file check) Fails until LATEST is fixed by hand No duplicate No duplicate No duplicate
scan.mode=latest start Stale id Stale id Real next snapshot Real next snapshot
incremental-between-timestamp Fails until LATEST is fixed by hand Reads all retained snapshots (5 rows instead of 2) Correct Correct
rollbackTo Returns normally, but keeps the snapshots after the target and deletes newer tags Same as master Correct Correct
Rollback procedures Fail at latestSnapshot() Reach rollbackTo above Correct Correct

2. A table that is already in the stale state

master A B A + B
New commit, committer replay Fail until LATEST is fixed by hand Succeed, no duplicate Same as master Succeed, no duplicate
incremental-between-timestamp Fails until LATEST is fixed by hand Reads all retained snapshots Same as master Reads all retained snapshots
scan.mode=latest start, rollbackTo Wrong already on master Same as master Same as master Same as master

With B in place, A only runs on tables that are already stale. There it brings back commits but turns other failures into silently wrong results, and it does not fix the rest. So I'd go with B alone.

Tests

The tests produce the stale hint through real commits with LATEST writes skipped, instead of resetting LATEST by hand. They cover the replay in both conflict-check modes, the Flink END_INPUT replay, scan.mode=latest, rollbackTo and incremental-between-timestamp. All of them fail on master and pass with B. The paimon-core and flink committer test suites pass as well.

Does this direction match how you see the hint files, and would you like the extra warning in the commit retry path as part of this PR? If so, I'll update the PR accordingly.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants