kotwal-itpro opened a new pull request, #29396:
URL: https://github.com/apache/flink/pull/29396

   ## What is the purpose of the change
   
   When a job is restored from the checkpoint whose id is set in 
`execution.state-recovery.without-channel-state.checkpoint-id`, 
`CheckpointCoordinator#extractOperatorStates` replaces every `OperatorState` 
with `OperatorState#copyAndDiscardInFlightData()`. That copy only carries over 
the subtask states and drops the coordinator state, so 
`restoreStateToCoordinators` resets every operator coordinator with `null`.
   
   For FLIP-27 sources this means `Source#createEnumerator` is called instead 
of `Source#restoreEnumerator`, and the restored job reads data again that it 
had already emitted (the JIRA has a `NumberSequenceSource` reproduction that 
emits the whole sequence a second time). The option is meant to drop in-flight 
data only, which lives in the subtask state, so the coordinator state should be 
kept.
   
   ## Brief change log
   
     - `OperatorState#copyAndDiscardInFlightData()` copies the coordinator 
state into the new `OperatorState`.
     - `MockOperatorCoordinatorCheckpointContext` records the data it is reset 
with, so tests can assert on it.
   
   ## Verifying this change
   
   This change added tests and can be verified as follows:
   
     - Added 
`CheckpointCoordinatorRestoringTest#testRestoreCoordinatorStateWithoutInFlightData`,
 which restores a checkpoint with coordinator state and input channel state 
while the in-flight data of that checkpoint is ignored. It checks that the 
coordinator is reset with the checkpointed bytes and that the input channel 
state is still dropped. Without the fix the coordinator receives `null` and the 
test fails.
     - `CheckpointCoordinatorTest`, `CheckpointCoordinatorRestoringTest`, 
`CheckpointCoordinatorTriggeringTest` and `OperatorCoordinatorHolderTest` pass.
   
   ## Does this pull request potentially affect one of the following parts:
   
     - Dependencies (does it add or upgrade a dependency): no
     - The public API, i.e., is any changed class annotated with 
`@Public(Evolving)`: no
     - The serializers: no
     - The runtime per-record code paths (performance sensitive): no
     - Anything that affects deployment or recovery: JobManager (and its 
components), Checkpointing, Kubernetes/Yarn, ZooKeeper: yes (restoring with 
`execution.state-recovery.without-channel-state.checkpoint-id` now restores 
operator coordinator state)
     - The S3 file system connector: no
   
   ## Documentation
   
     - Does this pull request introduce a new feature? no
     - If yes, how is the feature documented? not applicable
   
   ---
   
   ##### Was generative AI tooling used to co-author this PR?
   
   - [X] Yes (please specify the tool below)
   
   Generated-by: Claude Opus 5.5
   


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