88fantasy opened a new pull request, #12388:
URL: https://github.com/apache/seatunnel/pull/12388

   ### Purpose of this pull request
   
   Fixes #12386.
   
   After a master switch, `SubPlan#restorePipelineState` reports a RUNNING 
pipeline as `alreadyStarted` only if **every** physical vertex is `RUNNING`. 
For an unbounded source such as MySQL-CDC with `parallelism > 1`, idle readers 
are closed after the snapshot phase (`CloseIdleReaderOperation`), so their task 
groups are `FINISHED`. The restore therefore reports `alreadyStarted = false`, 
`CheckpointCoordinator#restoreCoordinator(false)` waits for `READY_START` from 
tasks that are already running (and will never report it again), and no 
checkpoint is ever triggered after the takeover. The job stays RUNNING, keeps 
reading and writing, but 2PC sinks never commit.
   
   Only flipping `alreadyStarted` is not enough: `closedIdleTask` is in-memory 
state of the old coordinator and is used to exclude closed subtasks from the 
ACK set (`getNotYetAcknowledgedTasks`), from barrier dispatch 
(`barrier.closedTasks()` in the enumerator) and from the committer's expected 
barrier count. Without it the new coordinator would wait for ACKs from closed 
readers.
   
   This PR:
   
   1. `SubPlan#restorePipelineState`: in a RUNNING pipeline, a `FINISHED` 
physical vertex no longer marks the pipeline as not started (bounded readers 
only finish after the final checkpoint, so here it can only be an idle-closed 
reader). These task groups are collected; coordinator vertices must still be 
`RUNNING`.
   2. The task groups are passed through 
`CheckpointManager#reportedPipelineRunning(int, boolean, 
Set<TaskGroupLocation>)` to `CheckpointCoordinator#restoreCoordinator(boolean, 
Set<TaskGroupLocation>)`, which rebuilds `closedIdleTask` from the plan after 
`cleanPendingCheckpoint` and before the `alreadyStarted` branch.
   3. The existing two-argument methods are kept and delegate with an empty 
set, so other callers (e.g. pipeline restart) are unchanged.
   
   No new persisted state is added; it relies only on the vertex states already 
restored from the IMap.
   
   ### Does this PR introduce _any_ user-facing change?
   
   Yes. With multiple masters, a CDC job with `parallelism > 1` keeps 
checkpointing after a master takeover instead of silently stalling.
   
   ### How was this patch tested?
   
   - Unit tests:
     - new 
`SubPlanRestoreStateTest#testRestoreRunningPipelineWithIdleClosedTaskGroupIsAlreadyStarted`
 (enumerator RUNNING, one reader group RUNNING, one FINISHED -> 
`reportedPipelineRunning(1, true, {group 3})`; before the fix it was `(1, 
false)`);
     - 
`CheckpointCoordinatorTest#testRestoreCoordinatorRebuildsClosedIdleTasksAfterMasterSwitch`
 (`closedIdleTask` equals the reader and writer of the closed group and a 
checkpoint is triggered).
     - `CheckpointCoordinatorTest`, `CheckpointManagerTest`, 
`SubPlanRestoreStateTest`, `StateTransitionCleanupTest`: 27/27 pass; reverting 
either part of the fix makes the new tests fail.
   - Manual test on a separated cluster (2 masters + 3 workers, JDK 8, 
checkpoint + IMap on S3), MySQL-CDC -> Doris (2PC, delete enabled), 
`parallelism = 2`, with the fix of #12387 also applied, after one source 
subtask became FINISHED:
     - `kill -9` the active master, run `DELETE`/`INSERT`/`UPDATE` on the 
source during the takeover: the new master logs `alreadyStarted: true, closed 
idle task groups: [TaskGroupLocation{..., taskGroupId=3}]`, checkpoints 
continue immediately, `SinkCommittedCount == SinkWriteCount`, no re-snapshot, 
sink row count and `SUM` equal to the source. Before the fix no checkpoint was 
triggered after the takeover and `SinkCommittedCount` stayed constant.
     - stop both masters, change the source while the cluster is down, restart 
everything: the pipeline is redeployed, `closedIdleTask` is cleared by the 
pipeline restart (`restoreCoordinator(false, [])`), checkpoints resume and the 
sink equals the source.
   
   Not covered: a master switch that happens exactly while an idle reader is 
being closed (in `readyToCloseIdleTask`, not yet acknowledged); behaviour there 
is unchanged by this PR.
   
   ### Check list
   
   * [ ] If any new Jar binary package adding in your PR, please add License 
Notice according
     [New License 
Guide](https://github.com/apache/seatunnel/blob/dev/docs/en/contribution/new-license.md)
   * [ ] If necessary, please update the documentation to describe the new 
feature. https://github.com/apache/seatunnel/tree/dev/docs
   * [ ] If necessary, please update `incompatible-changes.md` to describe the 
incompatibility caused by this PR.
   * [ ] If you are contributing the connector code, please check that the 
following files are updated:
     1. Update 
[plugin-mapping.properties](https://github.com/apache/seatunnel/blob/dev/plugin-mapping.properties)
 and add new connector information in it
     2. Update the pom file of 
[seatunnel-dist](https://github.com/apache/seatunnel/blob/dev/seatunnel-dist/pom.xml)
     3. Add ci label in 
[label-scope-conf](https://github.com/apache/seatunnel/blob/dev/.github/workflows/labeler/label-scope-conf.yml)
     4. Add e2e testcase in 
[seatunnel-e2e](https://github.com/apache/seatunnel/tree/dev/seatunnel-e2e/seatunnel-connector-v2-e2e/)
     5. Update connector 
[plugin_config](https://github.com/apache/seatunnel/blob/dev/config/plugin_config)
   


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