DanielLeens opened a new pull request, #12311:
URL: https://github.com/apache/seatunnel/pull/12311

   ## Purpose
   
   Fixes a real, recurring `engine-v2-it (8, ubuntu-latest)` failure that 
currently hits `dev` itself and many unrelated PRs:
   
   ```
   
SplitClusterFaultToleranceIT.testStreamJobCancelResolvesWhenWorkerCrashesBeforeCancelAck:448
     -> assertEventuallyCanceled:556
     org.opentest4j.AssertionFailedError: expected: <CANCELED> but was: <FAILED>
   ```
   
   On `dev`'s own push CI this leg has failed on every completed run since 
2026-09-13 04:01Z (5 in a row: runs 34736904718, 34747735985, 34752507232, 
34761938725, 34801745423) and intermittently before that since the test landed 
in #12030 — roughly a 70% failure rate, with the JDK 11 leg usually passing. 
The test is right about the intended semantics; the engine does not currently 
guarantee them in one specific window.
   
   ## Root cause
   
   A user cancels a streaming job while its worker dies. 
`PhysicalVertex#noticeTaskExecutionServiceCancel` sends `CancelTaskOperation` 
to the worker and, if the member leaves *before* acking, falls back to marking 
the vertex `CANCELED` locally (the #10729 fix). But the ack is immediate: 
`TaskExecutionService#cancelTaskGroup` just calls 
`cancellationFuture.cancel(false)` and returns, and the vertex's terminal state 
only arrives later through the task's own completion callback 
(`notifyTaskStatusToMaster`). So if the worker dies *after* acking but before 
that callback:
   
   1. `noticeTaskExecutionServiceCancel` has already returned normally at the 
`.invoke().get()` success path, leaving the vertex `CANCELING` and waiting for 
a callback that will never come (worker log: `Interrupted task 1000100000000` … 
then `notifyTaskStatusToMaster` → `NullPointerException: Target cannot be 
null!` because the worker's own cluster service is already gone).
   2. Hazelcast fires `memberRemoved` → 
`CoordinatorService#failedTaskOnMemberRemoved` → `makeTasksFailed`, which 
explicitly matches `CANCELING` vertices and marked them `FAILED` with "The 
taskGroup(...) deployed node(...) offline".
   3. `SubPlan#getPipelineEndState` sees `failedTaskNum > 0` → pipeline 
`FAILED` → `PhysicalPlan` → job `FAILED`.
   
   Which side of the ack the worker shutdown lands on is pure scheduling 
timing, so the same test passes or fails run to run. Verified end to end 
against `dev` run 34801745423's job log: the enumerator vertex (taskGroupId=1) 
went `FAILED` through exactly this `failedTaskOnMemberRemoved` stack, while the 
two reader vertices that were still inside the un-acked retry loop resolved 
`CANCELED` via the #10729 fallback and then rejected the later `FAILED` attempt 
("Task is trying to leave terminal state CANCELED").
   
   The test's Javadoc argued `failedTaskOnMemberRemoved` "cannot be the one to 
save this race" because the cancel chain holds the vertex lock; that is true 
only for the un-acked window and was the gap.
   
   ## Fix
   
   `CoordinatorService#makeTasksFailed` now resolves a lost vertex that is 
already `CANCELING` to `CANCELED` instead of `FAILED`, keeping the same "node 
offline" message for diagnostics. `DEPLOYING`/`RUNNING` vertices still resolve 
to `FAILED` (their work was genuinely cut short); all other states are left 
untouched, exactly as before. The decision is extracted into a small 
package-private `resolveLostMemberState(ExecutionState)` so it is unit-testable 
without a cluster.
   
   Why this is the right place rather than the test: a cancel was requested for 
the vertex and the worker acknowledged it; the member loss only removed the 
terminal callback. Resolving it as `CANCELED` keeps the outcome the cancel 
request was going to produce anyway, and matches what a clean cancel yields. 
Marking it `FAILED` reported a user-cancelled job as failed purely because a 
worker went away while honoring the cancel.
   
   **Scope of the behavior change** (checked against the callers):
   - User cancel + worker loss mid-cancel: job ends `CANCELED` (was: `FAILED`). 
This is the only observable change.
   - Restore is unaffected: a job in `CANCELING` has already called 
`JobMaster#neverNeedRestore()` (`PhysicalPlan#stateProcess` CANCELING/FAILING 
branch), and a `FAILING` pipeline that cancels its sibling tasks already has 
`failedTaskNum > 0` from the original failure, so it still ends `FAILED`.
   - No checkpoint, serialization, config, or API surface is touched.
   
   The E2E test's Javadoc is corrected to describe both mechanisms (un-acked → 
#10729 fallback; acked-then-lost → this path); its assertion is unchanged and 
remains the integration guard.
   
   ## Tests
   
   - New `CoordinatorServiceLostMemberResolutionTest` (pure unit test, no 
Hazelcast): `CANCELING → CANCELED`, `DEPLOYING`/`RUNNING → FAILED`, every other 
state (and `null`) → untouched.
   - Existing 
`SplitClusterFaultToleranceIT#testStreamJobCancelResolvesWhenWorkerCrashesBeforeCancelAck`
 becomes deterministic for the scenario it targets; its assertion is not 
relaxed.
   
   Per this repository's contribution guidance for Apache SeaTunnel, local 
execution was limited to `spotless:apply`/`spotless:check`; compilation and the 
tests above are verified by this PR's own GitHub Actions CI.
   
   🤖 Generated with [Claude Code](https://claude.com/claude-code)
   


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