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]