DanielLeens commented on code in PR #11841:
URL: https://github.com/apache/seatunnel/pull/11841#discussion_r3919904823
##########
seatunnel-engine/seatunnel-engine-core/src/main/java/org/apache/seatunnel/engine/core/job/JobImmutableInformation.java:
##########
@@ -167,6 +168,14 @@ public boolean isRestoreJob() {
return restoreMode != null && restoreMode.isRestore();
}
+ public boolean isSavepointRestore() {
Review Comment:
Good question to raise on a checkpoint/savepoint-adjacent field, but I don't
think the premise holds once you trace where `JobImmutableInformation` actually
gets serialized.
`JobImmutableInformation` itself is never embedded in checkpoint/savepoint
**state** storage (`PipelineState`/`CheckpointStorage`) — I grepped
`seatunnel-engine-server/.../checkpoint/**` and there's no reference to it
there. It's freshly *constructed* by the client at each resubmission, with
`restoreMode` explicitly set to `SAVEPOINT`/`CHECKPOINT` by the current-version
code doing the restoring
(`ClientJobExecutionEnvironment`/`RestJobExecutionEnvironment`). So for the
scenario you describe — user upgrades SeaTunnel, then restores from a savepoint
taken earlier — the `JobImmutableInformation` used for that restore submission
is built fresh by the new binary; it isn't deserialized from old bytes at all.
The two places it genuinely round-trips through serialization are the
RPC/REST submission wire and `runningJobInfoIMap` (used for master-node
failover recovery within a live cluster — see the comment at
`CoordinatorService.java:147`). Both go through `readData()`
(`JobImmutableInformation.java:243-269`), which is **untouched by this PR**
(verified via `git diff apache/dev...HEAD` on this file — the only changes are
the two new accessors at lines 171-177 and a one-line Javadoc). That method's
pre-existing fallback logic is:
```java
restoreMode = isStartWithSavePoint ? RestoreMode.SAVEPOINT :
RestoreMode.NONE; // line 259
restoreSourceJobId = isStartWithSavePoint ? jobId : null;
// line 260
if (hasRemainingBytes(in)) { ... }
// trailer, if present
```
This derives `restoreMode` from the legacy boolean *before* checking for the
trailer, so an old-format payload (pre-`RestoreMode`, pre-#11421) that was
genuinely savepoint-restored has `isStartWithSavePoint=true` on the wire and
correctly falls back to `RestoreMode.SAVEPOINT` — not "all false."
`isSavepointRestore()`/`isCheckpointRestore()` are pure derivations of that
already-correct `restoreMode` value, so they inherit the same backward-compat
guarantee without adding any new deserialization logic themselves.
So: no bug here that I can find — the compatibility path you're worried
about was already handled before this PR touched the file, and this PR's own
diff to this class is additive-only (two accessor methods + a comment). If you
have a specific call path in mind where `JobImmutableInformation` *does* get
persisted to durable storage and read back post-upgrade (rather than freshly
built per submission), I'd genuinely like to see it — that would be a real
pre-existing gap worth its own issue, just not one introduced by this diff.
##########
seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/master/JobMaster.java:
##########
@@ -284,7 +285,7 @@ public synchronized void init(long initializationTimestamp,
boolean restart) thr
sink.getId()));
JobMaster.handleSaveMode(
((SinkAction<?, ?, ?, ?>)
sink).getSink(),
- logicalDag.isStartWithSavePoint());
Review Comment:
Traced this through `LogicalDag`/`LogicalDagGenerator` to answer precisely
rather than just from the diff.
`LogicalDag.isStartWithSavePoint` (`LogicalDag.java:63`) is **untouched** by
this PR — it's still a plain `boolean` field. It's populated at DAG-generation
time via `AbstractJobEnvironment.buildLogicalDagGenerator` → `new
LogicalDagGenerator(actions, jobConfig, idGenerator, restoreMode.isRestore())`
(`AbstractJobEnvironment.java:135`, also untouched), which is already the
**broad** "any restore" truth value (true for both `SAVEPOINT` and
`CHECKPOINT`), despite the field's savepoint-sounding name. So
`logicalDag.isStartWithSavePoint()` and
`jobImmutableInformation.isRestoreJob()`/`getRestoreMode().isRestore()` were
already computing the identical boolean for every reachable path before this PR
— the swap at `JobMaster.java:273` and `:287-288` is value-preserving, not a
behavior change.
On your actual question — should `RestoreMode` live on `LogicalDag` instead
of `JobMaster` reaching into `JobImmutableInformation` directly: I'd argue the
current shape is actually the *more* honest one, not the ambiguous one.
`LogicalDag`'s field is a `boolean`, so if we wanted `LogicalDag` to
distinguish `SAVEPOINT` vs `CHECKPOINT` for any future use, it would need to
become a `RestoreMode`-typed field anyway — which would just be a second,
derived copy of the exact same value `JobImmutableInformation` already carries
authoritatively. `JobMaster` already holds a direct reference to
`jobImmutableInformation` (it's a constructor field, see
`jobImmutableInformation.getJobConfig()` two lines above at
`JobMaster.java:242`) and is the object that originally caused `logicalDag`'s
own flag to be set in the first place (client-side, via the same
`restoreMode`). So going straight to `jobImmutableInformation.getRestoreMode()`
at the master-side call site removes a layer of indirection
through a misleadingly-named boolean proxy, rather than introducing a new
coupling — `JobMaster` isn't reaching into someone else's internals here, it's
using the source of truth it already owns.
I don't see a separation-of-concerns issue in practice, but I take the point
that `LogicalDag.isStartWithSavePoint`'s name is now doubly misleading (it was
already broad-but-narrowly-named before this PR, and this PR's own
`CheckpointCoordinator` parameter rename — `isStartWithSavePoint` →
`isRestoreJob` — fixes exactly this class of naming problem elsewhere in the
same diff). Renaming
`LogicalDag.isStartWithSavePoint()`/`setStartWithSavePoint()` to something like
`isRestoreJob()` would be a reasonable, low-risk follow-up for consistency,
though it's a wire-serialized field (`LogicalDag.java:158,174`) so it'd need
the same "keep the old accessor name, just clarify semantics via Javadoc"
treatment this PR already used on
`JobImmutableInformation.isStartWithSavePoint()` rather than an outright
rename, to avoid a gratuitous compat-sensitive touch. Happy to see that filed
separately if you'd like to pursue it — not a blocker for this PR either way.
##########
seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/CoordinatorService.java:
##########
@@ -1394,9 +1396,12 @@ public PassiveCompletableFuture<Void> submitJob(
try {
JobImmutableInformation
submittedJobImmutableInformation =
deserializeJobImmutableInformation(jobImmutableInformation);
+ boolean isSavepointRestore =
Review Comment:
This is a fair, sharp catch — I checked and you're right that it's dead
inside this method: `submitJob(long jobId, Data jobImmutableInformation,
boolean isStartWithSavePoint)` (`CoordinatorService.java:1374-1375`) never
references the `isStartWithSavePoint` parameter anywhere in its ~90-line body
once you get past the signature — I grepped the full method and the only
occurrence is the declaration itself. Everything downstream uses
`submittedJobImmutableInformation.isSavepointRestore()`, derived from the
deserialized payload (line 1399), exactly as you traced.
Two things worth putting on the record here, though:
1. **This exact tension was already raised and resolved in this PR's own
thread**, before your comments — see @SEZ9's question and my answer in the
top-level conversation on 2026-08-20/21 (issue comments). The explicit,
deliberate decision was: keep the legacy wire argument as-is in this PR for
compatibility, and treat any deprecation/signature cleanup of the now-unused
parameter as a dedicated follow-up once every reader has migrated off it,
rather than bundling a signature change into an already-reviewed, additive
refactor. That's consistent with what you're describing — the parameter
*should* eventually go away or be replaced with an explicit `RestoreMode`, just
not in this diff.
2. **It's not purely inert further downstream, which is why I wouldn't call
it "obsolete" outright** — the same boolean is also what gets written onto
`SubmitJobOperation`'s own wire format (`writeInternal`/`readInternal`,
untouched by this PR — see your other comment on `BaseService.java:1360`, I've
answered the wire-level angle there). That's a genuinely separate concern from
`CoordinatorService.submitJob`'s *internal* decision-making: during a rolling
upgrade, an old-version master node doesn't know about
`JobImmutableInformation.isSavepointRestore()`/`RestoreMode` at all and needs a
coherent `isStartWithSavePoint` value to still be present on the wire it can
read. So the parameter has to keep existing at the RPC boundary even after it
becomes irrelevant to this method's own logic — which is exactly the "kept for
wire compat, no longer trusted for decisions" framing in my Compatibility
Impact section of the review above.
So: agreed this deserves a real cleanup, agreed it isn't currently reflected
in the method signature, but I'd treat it as the deprecation follow-up already
discussed with @SEZ9 rather than something this PR needs to absorb — happy to
see a tracking issue opened for "thread `RestoreMode` explicitly through
`CoordinatorService.submitJob`/`SubmitJobOperation` and retire the legacy
boolean" if you or @JeremyXin want to file it.
##########
seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/service/BaseService.java:
##########
@@ -1360,7 +1360,7 @@ protected JsonObject submitJobInternal(
new SubmitJobOperation(
Review Comment:
Same underlying observation as your comment on
`CoordinatorService.java:1399`, traced through the Hazelcast operation layer
instead — I answered the "should this be threaded explicitly / is it now dead"
question there. On the "silent override" framing specifically, though, I don't
think there's an actual override risk today, only redundancy:
At both of `BaseService.java`'s own producers of this value — the cross-node
RPC path at line 1360 you quoted, and the same-node direct path at line 1397
(`submitJob(...)` private helper) — the wire boolean passed into
`SubmitJobOperation`/`CoordinatorService.submitJob` is
`jobImmutableInformation.isSavepointRestore()`. That's the *same*
`RestoreMode`-derived value that `CoordinatorService.submitJob` later
re-derives from the deserialized payload
(`submittedJobImmutableInformation.isSavepointRestore()`,
`CoordinatorService.java:1399`). Same is true of `ClientJobProxy.java:77`, the
third producer. So for every currently-reachable caller in this codebase, the
value written onto the wire and the value the coordinator re-derives are
byte-for-byte identical by construction — there's no case in the current diff
where the outer parameter says one thing and the payload says another.
The scenario where they *could* genuinely diverge is a caller that predates
this consolidation (a not-yet-migrated or old-version producer sending a
stale/inconsistent `isStartWithSavePoint` alongside a payload with a different
real `RestoreMode`) — and that's precisely what the new
`CoordinatorServiceJobCleanupTest.testSubmitCheckpointRestoreDoesNotTrustLegacySavepointFlag`
test simulates and asserts against (submits `RestoreMode.CHECKPOINT` with the
wire boolean forced to `true`, asserts the checkpoint path — not savepoint's
cleanup-ownership branch — is what runs). So `CoordinatorService` deliberately
no longer trusts the RPC-level boolean for exactly the reason you're pointing
at: it *could* become stale/inconsistent from a future or mixed-version caller,
even though every current in-repo caller keeps it consistent. That's a
hardening, not a bug being papered over.
`SubmitJobOperation.writeInternal`/`readInternal` (the code you pasted) is
untouched by this PR and needs to keep writing/reading that boolean regardless,
for rolling-upgrade wire compatibility with an old-version node on the other
end — see my fuller answer on the `CoordinatorService.java:1399` thread for why
it can't just be dropped from the wire format even though it's no longer
trusted for decisions on the receiving side.
--
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]