[FLINK-40209][checkpoint] Introduce Regional Checkpoint core mechanism and notification dispatch - #28963
Open
raoraoxiong wants to merge 5 commits into
Open
Conversation
Collaborator
…iguration and refCheckpointId infrastructure - Add CheckpointListener.notifyRegionalCheckpointComplete(long, RegionalCheckpointInfo) for healthy-region tasks - Add CheckpointListener.notifyRegionalCheckpointFallback(long, long) for failed-region tasks - Add RegionalCheckpointInfo class with fallback checkpoint subtask mapping - Add OperatorCoordinator.supportsRegionCheckpoint() and checkpointCoordinatorForRegionFallback() - Wire OperatorCoordinatorCheckpointContext and OperatorCoordinatorHolder for forwarding - Add 3 config options: region.enabled, region.max-failure-ratio, region.max-consecutive-failures - Add refCheckpointId field to OperatorSubtaskState for tracking historical state references - Add MetadataV7Serializer for backward-compatible serialization of refCheckpointId - Add refCheckpointId to SubtaskStateStats/TaskStateStats for REST API aggregation - Add CheckpointSubsumeHelper for reference-aware checkpoint cleanup - Add regional config fields to CheckpointCoordinatorConfiguration Generated-by: CodeBuddy Code (GLM-5.2)
…d SourceCoordinator fallback - Add RegionalCheckpointHandler: decline buffering, region failure judgment, state recombination, two-tier max-consecutive-failures (Tier 1: force next global; Tier 2: abort + reset) - Wire CheckpointCoordinator to delegate regional checkpoint logic to RegionalCheckpointHandler - Add PendingCheckpoint methods: recordDecline, areAllTasksResponded, markUnacknowledgedTasksAsDeclined, reportFallbackSubtaskStats, finalizeRegionalCheckpoint - Add DefaultCompletedCheckpointStore.computeReferencedCheckpointIds for cleaner reference protection - Implement SourceCoordinator.supportsRegionCheckpoint() and checkpointCoordinatorForRegionFallback() - Implement SplitAssignmentTracker per-checkpoint assignment history with backward-compatible serialization - Wire Regional Checkpoint config through StreamGraph - Add unit tests: success path, consecutive limit, state assembly, deferred abort, cleaner, config Generated-by: CodeBuddy Code (GLM-5.2)
raoraoxiong
force-pushed
the
raorao/FLINK-40209-regional-checkpoint-core
branch
from
August 17, 2026 09:07
d79bacb to
d73109a
Compare
…dispatch and local state cleanup - Extend confirmCheckpoint RPC with fallbackCheckpointId parameter (reuses task-side checkpoint-complete RPC path so notification survives task restarts) - Add Task.notifyRegionalCheckpointFallback + CheckpointableTask.notifyRegionalCheckpointFallbackAsync - Implement StreamTask.notifyRegionalCheckpointFallbackAsync with SubtaskCheckpointCoordinator - Add SubtaskCheckpointCoordinator.notifyRegionalCheckpointFallback + OperatorChain propagation - Add AbstractUdfStreamOperator forwarding to user functions - Add TaskStateManager.pruneStateForCheckpoint for local state cleanup (FLIP-600 Section 9) - TaskExecutor.confirmCheckpoint dispatches to notifyRegionalCheckpointFallback or notifyCheckpointComplete Generated-by: CodeBuddy Code (GLM-5.2)
raoraoxiong
force-pushed
the
raorao/FLINK-40209-regional-checkpoint-core
branch
from
August 17, 2026 11:54
d73109a to
3bbabd8
Compare
…@PublicEvolving The regional checkpoint notification methods added to the @public CheckpointListener interface broke two CI checks: - ArchUnit PUBLIC_API_METHODS_USE_ONLY_PUBLIC_API_TYPES failed because notifyRegionalCheckpointComplete takes a @PublicEvolving RegionalCheckpointInfo argument, while public methods of a @public type may only expose @public leaf types. Annotating the method itself as @PublicEvolving moves it under the @PublicEvolving rule, which permits @PublicEvolving argument types. notifyRegionalCheckpointFallback gets the same annotation for consistency, since both are new FLIP-600 API that should not yet carry @public stability guarantees. - AbstractUdfStreamOperatorLifecycleTest#testAllMethodsRegisteredInTest failed because StreamOperator extends CheckpointListener, so both new default methods surface in StreamOperator.class.getMethods() and must be registered in the expected method list. Generated-by: CodeBuddy Code
…int options ConfigOptionsDocsCompletenessITCase failed because the three regional checkpoint options added to CheckpointingOptions were not present in the generated documentation: execution.checkpointing.region.enabled execution.checkpointing.region.max-consecutive-failures execution.checkpointing.region.max-failure-ratio The options already carry @Documentation.Section(EXPERT_CHECKPOINTING), so only the generated HTML was missing. Regenerated with: mvn package -Dgenerate-config-docs -pl flink-docs -am -nsu -DskipTests Generated-by: CodeBuddy Code
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
Implements the core Regional Checkpoint mechanism (FLIP-600 Phase 1-5): when partial pipeline regions fail during a checkpoint, the framework generates a logically complete Completed Checkpoint by combining historical state of failed regions with current state of healthy regions.
Changes
Commit 1: Define interfaces, configuration and refCheckpointId infrastructure
notifyRegionalCheckpointComplete(healthy-region tasks) andnotifyRegionalCheckpointFallback(failed-region tasks)supportsRegionCheckpoint()andcheckpointCoordinatorForRegionFallback()region.enabled,region.max-failure-ratio,region.max-consecutive-failuresrefCheckpointIdfield to OperatorSubtaskState + MetadataV7Serializer (backward-compatible)Commit 2: Implement core logic and SourceCoordinator fallback
Commit 3: Implement notification dispatch and local state cleanup
fallbackCheckpointIdparameter (reuses task-side checkpoint-complete RPC path so notification survives task restarts)Testing
Dependencies
Generated-by: CodeBuddy Code (GLM-5.2)