developer-rpai opened a new pull request, #18270: URL: https://github.com/apache/iceberg/pull/18270
## Problem After a **stateless restart** (no checkpoint/savepoint restored, `isRestored() == false`) of a Flink job using the modern `IcebergSink` (Sink V2 API), the job keeps consuming records and checkpoints succeed, but **no new Iceberg snapshots are ever committed again**. Data/manifest files keep being written to storage, but snapshot history stays frozen — with no errors or warnings. This is silent data loss. Fixes #18098. ## Root cause `IcebergCommitter.commit()` unconditionally calls `SinkUtil.getMaxCommittedCheckpointId(table, jobId, operatorId, branch)`, which walks the snapshot chain backwards for a snapshot matching the `(job-id, operator-id)` pair. Since that pair stays stable across a stateless restart, the lookup finds a snapshot from the **previous run** and returns its old, high checkpoint id (e.g. `13304`). The new run's checkpoint counter restarts at `1`, so every new committable falls into `headMap(maxCommittedCheckpointId, true)` and is marked already-committed via `signalAlreadyCommitted()` — indefinitely, until the counter organically exceeds the old value. The commits are skipped, not retried. The legacy `FlinkSink`/`FlinkFilesCommitter` guards this exact lookup with `isRestored()` and skips it on a stateless start. The Sink V2 `IcebergCommitter` had no such guard. ## Fix - `IcebergSink.createCommitter()` now derives restored-ness from `CommitterInitContext.getRestoredCheckpointId().isPresent()` (the Sink V2 API's equivalent of `isRestored()`, available since Flink 1.18) and passes it into `IcebergCommitter`. - `IcebergCommitter.commit()` only consults the table's commit history when the job was restored; on a fresh start it begins from `SinkUtil.INITIAL_CHECKPOINT_ID` (-1), exactly like the legacy sink path, so all new commits proceed. - Applied identically to all four Flink version modules (`flink/v1.20`, `v2.1`, `v2.2`, `v2.3`). Exactly-once semantics on a **stateful** restore are unchanged: the lookup still runs and already-committed checkpoints are still skipped. ## Tests - Added `testStatelessRestartCommitsNewCheckpoints` to `TestIcebergCommitter` (all four modules): simulates a previous run committing checkpoints 1–5 with the same job/operator ids, then commits checkpoints 1–3 through a fresh (not-restored) committer. It asserts new snapshots are created (5 → 8), `signalAlreadyCommitted()` is never invoked, the max committed checkpoint id tracks the new run, and all rows are readable. - Existing `getCommitter()` call sites now pass `isRestored=true`, preserving current behavior for the restore-path tests (`testStateRestoreFromPreJobWithCommitted`, etc.). ## Verification level I traced the full code path on `main` (`IcebergSink.createCommitter` → `IcebergCommitter.commit` → `SinkUtil.getMaxCommittedCheckpointId`) and confirmed the bug is present as described in #18098, including the legacy `FlinkFilesCommitter` `isRestored()` guard this mirrors. I could not execute the Flink test suite in this environment (no Maven/Flink runtime available); the new/changed files were syntax-checked, and the added test follows the existing `TestIcebergCommitter` patterns. CI should run the full `TestIcebergCommitter` suite. ## Impact Fixes silent data loss for every `IcebergSink` user whose job restarts without restoring state (e.g. fresh redeploys on managed Flink services) while keeping the same job/operator ids — previously all post-restart writes were silently dropped until the checkpoint counter caught up with the old run's value. -- 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] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
