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]

Reply via email to