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]

Reply via email to