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]