[flink] Support unaligned checkpoints for coordinator commit - #9451
Merged
JingsongLi merged 1 commit intoAug 28, 2026
Merged
Conversation
leaves12138
approved these changes
Aug 28, 2026
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Purpose
FlinkSink currently rejects unaligned checkpoints for every streaming commit path. Coordinator commit performs the commit in the writer's OperatorCoordinator and can support unaligned checkpoints without relying on the classic global committer operator.
This PR allows unaligned checkpoints only when coordinator commit is enabled. The classic operator-commit path continues to reject unaligned checkpoints, and all paths still require EXACTLY_ONCE checkpointing.
The tests also cover the actual append-table use case with
data-evolution.enabled,row-tracking.enabled, bucket-unaware mode, and coordinator commit enabled.Tests
mvn -pl paimon-flink/paimon-flink-common -Pflink1 -DwildcardSuites=none -Dtest=FlinkSinkTest,CoordinatorCommitITCase clean test(Java 8)mvn -pl paimon-flink/paimon-flink-common -Pflink2 -DwildcardSuites=none -Dtest=FlinkSinkTest,CoordinatorCommitITCase test(Java 11)Both runs passed 15 tests with no failures.