developer-rpai commented on issue #18098: URL: https://github.com/apache/iceberg/issues/18098#issuecomment-5843734026
I've opened PR #18270 with a fix for this. **Root cause (confirmed on main):** `IcebergCommitter.commit()` (the Sink V2 committer used by `IcebergSink`) calls `SinkUtil.getMaxCommittedCheckpointId(table, jobId, operatorId, branch)` unconditionally. Since `(job-id, operator-id)` stay stable across a stateless restart, the lookup matches the previous run's snapshots and returns its high checkpoint id, so every new committable (checkpoint counter restarted at 1) is marked already-committed and never committed — silent data loss. **Fix:** gate the lookup on whether the job was restored, mirroring the legacy `FlinkSink`/`FlinkFilesCommitter` `isRestored()` guard. `IcebergSink.createCommitter()` now passes `context.getRestoredCheckpointId().isPresent()` into `IcebergCommitter`, which skips the history lookup on a fresh start (starting from `INITIAL_CHECKPOINT_ID`, like the legacy path). Applied to all four Flink modules (v1.20/v2.1/v2.2/v2.3), plus a regression test simulating the stateless-restart scenario. Happy to adjust the approach if maintainers prefer a different fix shape. -- 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]
