88fantasy opened a new issue, #12385:
URL: https://github.com/apache/seatunnel/issues/12385
### Search before asking
- [x] I searched the issues and found no similar issues.
### What happened
In a separated Zeta cluster with checkpoint storage and IMap persistence
enabled, a streaming job that is restored after a **master switch / master
restart** does not load its latest checkpoint. For MySQL-CDC this means the
source runs a **full snapshot again** and the sink starts from an empty state,
so every `DELETE` executed on the source while the master was down is lost in
the sink (INSERTs/UPDATEs are repaired by the snapshot, deletes are not).
### Root cause
This is a regression introduced by #11421 (`2b51d3614`).
`JobMaster#initCheckPointManager` passes `isRestoreJob() || restart`,
`getRestoreMode()` and `getRestoreSourceJobId()` to `CheckpointManager`. For a
normally submitted job `restoreSourceJobId` is `null` (also after it is
deserialized from the IMap), and since #11421 `CheckpointManager` only loads a
checkpoint when
```java
checkpointConfig.isCheckpointEnable() && isRestoreJob && restoreSourceJobId
!= null
```
so on `restart = true`
(`CoordinatorService#restoreJobFromMasterActiveSwitch` -> `jobMaster.init(...,
true)`) `pipelineState` stays `null`, `latestCompletedCheckpoint` is never set
and no `Restore checkpoint, job id: ...` is logged. Before #11421 the restart
path loaded `getLatestCheckpointByJobIdAndPipelineId(jobId, ...)` of the job
itself. The comment in `CheckpointCoordinator` ("For job restore from master
node active switch") shows this path is meant to restore state.
A second, related effect: a job that was started with `-r` /
`--restore-with-checkpoint` and later hits a master switch is restored from the
**original** savepoint / source job instead of its own latest checkpoint, i.e.
its state goes back in time.
### How to reproduce
1. Separated mode, checkpoint storage on S3 (hdfs plugin, `storage.type:
s3`), IMap map-store enabled, checkpoint interval 10s.
2. Submit a MySQL-CDC -> Doris (2PC, delete enabled) streaming job, wait for
the snapshot to finish and a few checkpoints to complete.
3. Stop all masters (SIGTERM). Workers shut down by themselves (all lite
members).
4. On the source: `DELETE` 5 rows, `INSERT` 20 rows, `UPDATE` 5 rows.
5. Start masters and workers again. The job is restored as RUNNING (`begin
restore job ... from master active switch`).
Actual: `MySqlSnapshotSplitReadTask - Snapshot step 1 ...` (full snapshot
again), `DorisSinkWriter - restore checkpointId 0`, no `Restore checkpoint` log
on the master. Sink ends with 5 more rows than the source (exactly the deleted
ones).
Expected: `pipeline(1) restore with ... on checkpointId(N)` / `Restore
checkpoint`, no new snapshot, sink equals source.
### Fix
I have a small patch: `JobMaster` passes the `restart` flag to
`CheckpointManager` (new constructor overload; the existing one delegates with
`false`). On master failover the job's own latest checkpoint is preferred,
falling back to `restoreSourceJobId`/`restoreMode` only when the job has no
checkpoint yet. Unit tests added to `CheckpointManagerTest` (own checkpoint on
failover, own checkpoint preferred over the original savepoint, fallback to
source job).
Verified on a 2-master / 3-worker cluster: after the patch, both "restart
all masters" and "worker failure after a master takeover" resume from the
latest checkpoint without re-snapshotting, and deletes/inserts/updates made
during the downtime are all applied (sink row count and SUM equal to the
source). Happy to open a PR with this change.
### SeaTunnel Version
dev (`c7304ace6`, 2026-09-17), regression since #11421
### SeaTunnel Config
```conf
env { job.mode = "STREAMING", parallelism = 2, checkpoint.interval = 10000 }
source { MySQL-CDC { table-names = ["db.orders"], startup.mode = "initial",
... } }
sink { Doris { sink.enable-2pc = "true", sink.enable-delete = "true", ... } }
```
### Running Command
```shell
bin/seatunnel.sh --config cdc.conf --async
```
### Error Exception
```log
No exception. On the restored master: "begin restore job (...) from master
active switch" but no "Restore checkpoint"; on the workers: "Snapshot step 1 -
Determining low watermark" and "restore checkpointId 0".
```
### Zeta or Flink or Spark Version
Zeta, separated cluster mode
### Java or Scala Version
JDK 8 (1.8.0)
### Screenshots
_No response_
### Are you willing to submit PR?
- [x] Yes I am willing to submit a PR!
### Code of Conduct
- [x] I agree to follow this project's Code of Conduct
--
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]