fishfishfishfishaa opened a new pull request, #8671:
URL: https://github.com/apache/paimon/pull/8671

   ## Purpose
   
   This PR implements the end-of-input handling that was intentionally left out 
of the core coordinator-commit PR for [PIP-30 
([[#8220](https://github.com/apache/paimon/issues/8220)](https://github.com/apache/paimon/issues/8220))].
   
   It extends `sink.coordinator-commit.enabled` to correctly finish and commit 
bounded inputs. In particular, coordinator commit now supports batch jobs 
without checkpointing, and correctly handles bounded streaming jobs after all 
writers reach `endInput`.
   
   This PR is built on top of the core coordinator-commit implementation. It 
does not change the default committer-operator path.
   
   ## What changes
   
   ### Commit after all writers finish
   
   Each coordinated writer emits a final committable entry with 
`Long.MAX_VALUE` as a dedicated end-of-input checkpoint ID when Flink invokes 
`endInput`.
   
   The JobManager-side `CommittingWriteOperatorCoordinator` tracks this entry 
separately from ordinary checkpoint committables:
   
   - An end-input entry covers all later ordinary checkpoints for that writer, 
because a finished writer cannot produce more data.
   - Once every writer has reported end input, the coordinator performs one 
final commit through `filterAndCommit(..., false, true)`.
   - After the final commit succeeds, repeated end-input events and later 
checkpoint-complete notifications are ignored.
   
   For streaming jobs with checkpointing enabled, the final commit remains 
aligned with checkpoint completion. This preserves the normal checkpoint-driven 
commit model while allowing some subtasks to finish before others.
   
   For batch jobs, where checkpointing is normally disabled, the coordinator 
commits immediately after receiving end-input entries from all writer subtasks.
   
   ### Support batch coordinator commit
   
   Coordinator commit previously required streaming mode with checkpointing 
enabled. This is now relaxed as follows:
   
   - Streaming mode still requires checkpointing.
   - Batch mode is supported without checkpointing and commits when all writers 
reach end input.
   
   The existing coordinator-commit restrictions remain unchanged: it only 
supports write-only unaware-bucket append tables, and still rejects 
configurations such as primary-key tables, precommit compaction, 
auto-tag-for-savepoint, and concurrent checkpoints.
   
   The option documentation is updated accordingly.
   
   ### Preserve end-input state across recovery
   
   End input is not an ordinary checkpoint: it is newer than every checkpoint 
and may need to survive restore even if Flink does not invoke `endInput` again 
after recovery.
   
   To make this safe:
   
   - Writers persist the final end-input committable in their independent 
operator state.
   - On restore, writers replay the persisted final entry to the coordinator.
   - If `endInput` is invoked again after restore, the writer merges newly 
produced final committables with the persisted entry and sends one 
authoritative final entry.
   - `WriterCommittables` treats end input separately from the maximum ordinary 
checkpoint, so it does not incorrectly constrain ordinary checkpoint alignment.
   - Replayed end-input entries are replaceable rather than treated as 
duplicate ordinary checkpoint reports.
   
   This also handles region failover while the coordinator is still waiting for 
the remaining writers to finish.
   
   ### End-input watermark handling
   
   The writer now uses `end-input.watermark`, when configured, for the final 
end-input committable.
   
   For ordinary checkpoints after a subtask has finished, that subtask 
contributes `Long.MAX_VALUE` to the watermark minimum. This makes a finished 
writer neutral and prevents it from incorrectly holding back the watermark of 
active writers.
   
   The final commit uses the configured end-input watermark as expected.
   
   ## How it works
   
   ```text
   writer subtask reaches endInput
           |
           v
   emit final CheckpointCommittables(Long.MAX_VALUE)
           |
           v
   persist final entry in writer operator state
           |
           v
   send CommittableEvent to JobManager coordinator
           |
           v
   all writer subtasks have reported end input?
           |
           +-- no  --> keep waiting; ended subtasks cover later normal 
checkpoints
           |
           +-- yes --> perform one final filterAndCommit(..., false, true)
   ```
   
   ## Tests
   
   This PR adds coverage for:
   
   - Batch coordinator commit without checkpointing.
   - Streaming coordinator commit with bounded input.
   - End-to-end final data commit after input completion.
   - End-to-end propagation of `end-input.watermark`.
   - Writer-side emission, persistence, restore, and re-emission of final 
committables.
   - Coordinator-side final-commit behavior when checkpointing is disabled.
   - Recovery and region-failover handling for end-input entries.
   - `WriterCommittables` semantics:
     - end input covers later ordinary checkpoints;
     - finished subtasks do not constrain later watermark aggregation;
     - restored state may contain an end-input entry newer than the restored 
checkpoint;
     - repeated end-input reports replace the authoritative final entry;
     - ordinary cleanup retains end-input state until the final commit.
   
   No production behavior changes when `sink.coordinator-commit.enabled` is 
disabled.


-- 
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]

Reply via email to