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

   ### Purpose of this pull request
   
   Fixes #12385.
   
   Since #11421 a job that is restored after a master switch / master restart 
(`JobMaster#init(..., restart = true)`) no longer loads any checkpoint, because 
`CheckpointManager` only loads state when `restoreSourceJobId != null`, which 
is always `null` for a normally submitted job. As a result MySQL-CDC sources 
re-run the full snapshot after a master failover and sinks start from an empty 
state, so deletes executed on the source while the master was down are lost.
   
   This PR passes the `restart` flag from `JobMaster` to `CheckpointManager` 
(new constructor overload; the existing constructor delegates with `false`, so 
other callers are unchanged). On master failover:
   
   1. the latest completed checkpoint of the job itself is loaded (the 
behaviour before #11421);
   2. only if the job has no checkpoint of its own yet, it falls back to 
`restoreSourceJobId` / `restoreMode` (e.g. a job started with 
`--restore-with-checkpoint` whose master died before its first checkpoint).
   
   This also fixes a job started with `-r` being restored from its original 
savepoint (instead of its latest checkpoint) when it later hits a master switch.
   
   Savepoint / `--restore-with-checkpoint` submissions (`restart = false`) keep 
the #11421 semantics.
   
   ### Does this PR introduce _any_ user-facing change?
   
   Yes. After a master switch/restart, running streaming jobs resume from their 
latest checkpoint again instead of restarting without state (behaviour before 
#11421).
   
   ### How was this patch tested?
   
   - Unit tests added in `CheckpointManagerTest`:
     - `testMasterFailoverRestoreShouldLoadLatestCheckpointOfOwnJob`
     - `testMasterFailoverRestoreShouldPreferOwnLatestCheckpointOverSavepoint`
     - `testMasterFailoverRestoreShouldFallBackToSourceJobWithoutOwnCheckpoint`
   
     `CheckpointManagerTest` 9/9 passes. With the failover branch disabled, the 
first two new tests fail (expected checkpoint id 8, actual 1 and 4).
   - Manual test on a separated cluster (2 masters + 3 workers, JDK 8, 
checkpoint + IMap on S3), MySQL-CDC -> Doris (2PC, delete enabled):
     - stop both masters, `DELETE` 5 / `INSERT` 20 / `UPDATE` 5 rows on the 
source, restart all nodes: master logs `pipeline(1) restore with NONE on 
checkpointId(24)` + `Restore checkpoint`, no new snapshot, sink row count and 
`SUM` equal to the source (before the patch: full re-snapshot, the 5 deleted 
rows remain in the sink);
     - after a master takeover, kill the worker running the tasks: pipeline 
restores from the checkpoint without re-snapshotting, sink equals source.
   
   Note: master takeover with CDC `parallelism > 1` still stops triggering 
checkpoints for a different reason, tracked in #12386 (not addressed here).
   
   ### 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