davidzollo opened a new pull request, #12029:
URL: https://github.com/apache/seatunnel/pull/12029
## What this guards
#10567 ("[Fix][Zeta] Make `deployTask` idempotent during master failover
recovery") fixed a bug
where `TaskExecutionService#deployTask` threw
`RuntimeException("TaskGroupLocation: ... already
exists")` whenever it was invoked for a `TaskGroupLocation` already present
in
`executionContexts`, i.e. a task that is genuinely still executing on this
worker. After the fix,
that call logs a warning, releases the classloaders it just acquired for the
redundant deploy
attempt, and returns `TaskDeployState.success()` as a no-op instead.
That collision is reachable through only one narrow window, confirmed by
reading the current
`PhysicalVertex.java`:
- On master failover, `initStateFuture()` restores each vertex's persisted
`ExecutionState` from
the running-job IMap and, for a persisted state of `RUNNING` or
`DEPLOYING`, pings the worker via
`checkTaskGroupIsExecuting` to self-heal a state that may no longer match
reality.
- For `RUNNING`, a confirmed-alive task is simply left as `RUNNING`, and
`SubPlan#stateProcess`'s
`case RUNNING` does nothing — no redeploy call is ever made, so
`RUNNING`-state tasks never hit
the bug.
- For `DEPLOYING`, the same confirmation still leaves the state as
`DEPLOYING` (the self-heal check
only ever downgrades to `FAILING`, it never upgrades to `RUNNING`), and
`stateProcess`'s `case
DEPLOYING` unconditionally calls `deploy()` again regardless of what the
check just found — this
is the exact pre-fix crash site.
- So the precise trigger is: kill the active master while at least one
task's persisted state is
`DEPLOYING`, and the underlying worker has (or will have, by the time
restore runs) already
succeeded in deploying it. `DEPLOYING` only persists in the IMap for the
duration of one worker
deploy RPC round trip (typically low single-digit milliseconds), so
hitting it reliably needs a
deliberately constructed trigger rather than a blind sleep.
## Coverage gap before this PR
The fix commit's own unit test,
`TaskExecutionServiceTest#testDeployTaskIdempotentWhenAlreadyRunning`,
calls `taskExecutionService.deployTask(data)` twice directly in the same
JVM. It never exercises
the real multi-node master-failover trigger path: no Hazelcast cluster, no
`initStateFuture`
self-heal, no `stateProcess` redeploy branch. I grepped the `seatunnel-e2e`
tree for `DEPLOYING`,
`checkTaskGroupIsExecuting`, and `deployTask` and found no existing E2E test
that touches any of
them. This PR closes that gap with a real 2-master/1-worker cluster.
## What the new test does
`SplitClusterFaultToleranceIT#testDeployNotDuplicatedWhenMasterKilledDuringDeploy`
submits a
bounded batch job (FakeSource → LocalFile,
`cluster_batch_fake_to_localfile_template.conf`, the
same template `ClusterFaultToleranceIT`/`SplitClusterFaultToleranceIT`
already use for exact
row-count restore assertions) with parallelism 40, and on a dedicated
watcher thread tight-polls
(1 ms between checks, no blind sleep) the active master's own in-JVM
`JobMaster#getPhysicalPlan()` for any task vertex reporting
`ExecutionState.DEPLOYING`. The instant
one is found, the watcher shuts the active master down right there —
reproducing a real master
crash inside the exact window #10567 guards. It then waits for the standby
to take over, asserts
the recovered job's every vertex is back to `RUNNING` with each pipeline's
`SubPlan#getPipelineRestoreNum()` at most 1, waits for the (unrestarted) job
to reach `FINISHED`,
and asserts the LocalFile output on the same "recovered and stable" terms
this file already uses
for its other master-down restore tests.
I placed this on `SplitClusterFaultToleranceIT` rather than
`SplitClusterPendingJobLifecycleFailoverIT`
(the class scenario I1's #12027 extended): that class's own Javadoc is "Test
the job recovery
capability ... in case of cluster node failure," it already has a
`testBatchJobRestoreInMasterDown`
sibling using the identical batch-template technique, and it uses split-role
master-only/worker-only nodes so the worker that must survive the kill is
never itself a kill
target — the closer thematic and infra fit for a "worker keeps running
through a master crash"
scenario than the pending-job-queue class.
## Why the trigger is reliable, not flaky
`SubPlan#stateProcess` dispatches deploys to a pipeline's vertices one at a
time on a single
thread: each vertex only becomes visibly `DEPLOYING` once its RPC to the
worker is in flight, and
stays that way until the ack returns. Raising parallelism therefore does not
widen any single
vertex's window — it multiplies how many independent windows the watcher
gets to land in across
the whole deploy fan-out (40 chances instead of 1). The watcher's own poll
is a plain in-JVM field
read (`PhysicalVertex#getExecutionState()`) with no network cost, while the
state it races against
is gated by a real IMap write plus a real deploy RPC round trip, so the
watcher can sample many
times inside a window it does not control. That asymmetry — cheap
high-frequency polling against a
comparatively expensive state transition, multiplied across many independent
vertices — is what
makes this deterministic rather than a hope-it-lands sleep. The test fails
loudly (rather than
silently passing a weaker test) if `DEPLOYING` is never observed within 30s,
so a future change
that breaks this assumption is caught explicitly instead of quietly
degrading into a no-op test.
Being honest about what the trigger does and doesn't guarantee: observing
`DEPLOYING` at kill time
guarantees the persisted state driving restore is `DEPLOYING` (necessary to
reach `stateProcess`'s
redeploy branch), but does not by itself guarantee the worker had already
finished deploying that
vertex (the other ingredient the pre-fix crash needs). Both sub-cases are
legitimate outcomes of
this trigger, and the assertions are written to hold cleanly on either one:
if the worker had
already succeeded, the fixed code makes the redeploy a no-op and that vertex
needs no restore at
all; if it had not, the same self-heal check correctly drives it to
`FAILING` and the pipeline
restores exactly once, which is ordinary, expected recovery rather than a
regression. What must
never happen on either sub-case is a client-visible failure or more than one
pipeline restore for
the whole scenario — a pre-fix run instead throws on the worker, fails that
one vertex, and (since
a single failed vertex fails its whole pipeline) forces a full pipeline
restart that redeploys
every other already-fine vertex in it as collateral damage, which is exactly
what the
`getPipelineRestoreNum() <= 1` assertion catches.
## Test plan
- [x] `./mvnw spotless:apply -pl
seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base -am -nsu` —
BUILD SUCCESS
- [x] `./mvnw install -pl
seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base -am -nsu
-Dmaven.gitcommitid.skip=true -DskipTests -Dmaven.test.skip=true
-Dspotless.check.skip=true -T 3C` — BUILD SUCCESS (confirms the new test
compiles and packages; `-DskipTests` per repo convention since the shaded
modules on this path only relocate classes during package/install, not bare
compile)
- [ ] Full E2E execution (`./mvnw -DskipUT -DskipIT=false verify` on this
module) is deferred to CI per this repository's local-verification conventions;
only the affected module was rebuilt locally, `seatunnel-engine-ui` was not
touched and was not rebuilt.
--
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]