corgy-w opened a new issue, #11961:
URL: https://github.com/apache/seatunnel/issues/11961

   ## Problem
   
   Zeta engine currently stores checkpoints and savepoints with the **same 
storage format** (runtime `PipelineState -> CompletedCheckpoint` blobs in the 
checkpoint directory, no version marker), which causes the following issues:
   
   1. **Savepoints can be silently lost**: job-terminal cleanup 
(`clearCheckpointIfNeed` / `cleanupTerminalZombieCheckpointIfNecessary`) 
deletes the whole `<namespace>/<job-id>/` directory, and a restored job reuses 
the same job id - so after `stop-with-savepoint -> upgrade -> restore -> 
finish/cancel`, the old savepoint is deleted together with the checkpoints. 
`max-retained` rotation can also delete savepoint files.
   2. **No format version**: the serialized payload has no version marker, so 
any structural change to the persisted classes is unguarded (ProtoStuff 
`RuntimeSchema` assigns field numbers by reflection order - 
inserting/reordering a field can silently mis-assign data instead of failing).
   3. **Documentation describes a layout that does not exist** (`savepoint/` 
directories + `state-data` files + `savepoint.path` env option were written in 
#10981 / #10979 as documentation-only; the code never implemented them).
   
   ## Proposed change
   
   Follow the industry-standard (Flink-like) design: **directory isolation + 
versioned metadata marker + versioned readers**.
   
   - **New engine-wire-v1 payload**: dedicated wire DTOs 
(`WireCheckpoint`/`WireActionState`/`WireSubtaskState`/`WireTaskStatistics`/`WireSubtaskStatistics`)
 with stable `@Tag` field numbers, name-based enum encoding, runtime-only 
fields (`isRestored`) excluded from the storage contract; `CheckpointWireCodec` 
converts between wire DTOs and runtime model; golden byte fixtures 
(`savepoint-wire/legacy-v0`) pin the legacy layout.
   - **Optional `SavepointStorage` capability**: new SPI in 
`checkpoint-storage-api` (begin/commit/abort writer, list/read/delete bundles, 
staging isolation, `SavepointMeta` + per-file manifest with SHA-256 checksums). 
The existing `CheckpointStorage` SPI stays unchanged; `LocalFileStorage` and 
`HdfsStorage` implement the capability.
   - **Engine integration**: `JobMaster.savePoint()` opens a single-flight 
`SavepointSaveSession`; coordinators write engine-wire payloads into 
`.staging/<attemptId>/`; the bundle is committed (publishing `_metadata.ser`) 
only after every pipeline reaches `SUSPEND`, otherwise aborted. Savepoints are 
excluded from checkpoint rotation and terminal cleanup and are kept until 
explicitly deleted.
   - **Versioned readers**: `SavepointReaderRegistry` dispatches by 
`formatVersion` (`v1` now); bundles below `MIN_SUPPORTED_FORMAT_VERSION` 
require migration, above the current window are rejected with an actionable 
`SavepointIncompatibleException`. Restore reads the newest completed bundle 
first and falls back to the legacy checkpoint-directory scan (best-effort) when 
no bundle exists.
   - **Docs**: `state-storage-and-recovery.md` (en/zh) corrected to the real 
layout, `restore.mode`/`savepoint.path` marked as planned-not-implemented, 
`incompatible-changes.md` records the change, legacy compatibility boundary and 
connector-state limitation.
   
   Scope notes: migration tool (option B), management APIs 
(list/delete/restore-by-id), multi-pipeline failure/failover protocol e2e and 
object-storage (S3/COS) integration checks are deliberately left as follow-ups 
documented in the design doc.
   
   Full design: 
[savepoint-checkpoint-storage-compat-design.md](https://github.com/apache/seatunnel/blob/dev/docs/design/savepoint-checkpoint-storage-compat-design.md)
   


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