fishfishfishfishaa commented on PR #8671: URL: https://github.com/apache/paimon/pull/8671#issuecomment-5834421042
Thanks for the detailed review and the follow-up question. You are right that restoring a checkpointed terminal batch is different from replaying input from an earlier checkpoint. The previous design did not make that distinction safely enough. I have reworked the protocol in my local revision. The main change is to commit each writer’s terminal tail under a real checkpoint ID and reserve `Long.MAX_VALUE` for data-free global finalization. #### 1. Restored terminal committables The replacement you described is unsafe: A belongs to the restored checkpoint, so its input is not replayed, and replacing it with an empty B could lose pending data. In the revised design, `endInput()` only seals the writer. The next checkpoint K captures the terminal tail and a terminal marker, persists them in writer state, and reports them through the checkpoint path. On recovery from K, the writer replays the checkpointed contribution and restores its sealed-input status. A repeated `endInput()` does not prepare a replacement batch. Subsequent checkpoints carry empty terminal markers without replacing the original checkpointed tail. The coordinator commits the tail through the ordinary checkpoint/recovery path. It records `terminalCoveredBy` only after successful reconciliation. Later coordinator checkpoints persist that coverage, so recovery no longer depends on an already-covered writer reporting again. #### 2. Failure propagation during completion Agreed. Draining `close()` only addresses unfinished work during normal shutdown. It does not establish that a commit/tag failure will reach the caller, and the normal non-drain stop-with-savepoint test cannot prove that. For EndInput completion, the revised design keeps a writer-side completion barrier. Early terminal writers may receive permission before their ordinary commit finishes, but they still wait for a completed Flink checkpoint covering their tail. The last terminal candidate remains waiting until ordinary commit, applicable tag work, and global MAX/listener finalization succeed. A coordinator failure fences subsequent work and prevents the final release. EndInput completion therefore does not rely on reporting failures from `close()`. This does not by itself resolve the separate non-drain stop-with-savepoint case, which cannot be assumed to execute the EndInput protocol. I will keep that limitation explicit rather than claim it is fixed by draining or by this final-release mechanism. #### 3. `end-input.watermark` The revised coordinator uses the configured `end-input.watermark` for the global MAX finalization committable. Ordinary checkpoint commits retain their normal aligned watermarks. This also means an early writer cannot apply the configured final watermark while another writer is still processing input. If the option is absent, finalization uses the last processed global watermark, retained in coordinator state. Supplying this watermark and creating a new snapshot are separate concerns; snapshot creation/filtering still follows the committer’s behavior. #### Why report from `writer.endInput()`? The original report was intended to transfer the final tail and tell the coordinator that this writer had no more input. It was not sufficient proof that the tail was covered by a successful checkpoint. The revised design removes that separate report from `endInput()`. The real-checkpoint contribution carries both the tail and the terminal marker, so EndInput does not introduce another commit trigger. The supported scope remains checkpoint-enabled streaming coordinator commit with one concurrent checkpoint. I have also written out the normal, abort, region/global recovery, and watermark scenarios here: [EndInput protocol and recovery scenarios](https://docs.google.com/document/d/15Pm5sIRRj6BxEKfuWdNO0p-26e65EsLy1cm49ikQQO4/edit?usp=sharing) I will update the PR description to describe this protocol and keep the validation results separate from the design reasoning. -- 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]
