DanielLeens opened a new pull request, #12199:
URL: https://github.com/apache/seatunnel/pull/12199
### Purpose of this pull request
Adds one new Zeta E2E regression test covering the fix in #10448 for #10442
("seatunnel will never perform a checkpoint again once a previous checkpoint
fails"), as part of
the ongoing extreme-case E2E coverage initiative for the Zeta engine.
**Issue #10442**: `CheckpointCoordinator#startTriggerPendingCheckpoint`
increments `pendingCounter`
before dispatching a checkpoint barrier. If that dispatch throws (the
reporter's example:
`CheckpointManager#sendOperationToMemberNode` ->
`JobMaster#queryTaskGroupAddress` timing out under
IMap contention), the old code just logged the exception and returned,
leaving `pendingCounter`
stuck above zero forever. Every later scheduled trigger attempt then saw
`pendingCounter > 0` and
just rescheduled itself, so the job kept reporting `RUNNING` with no error
and no further
checkpoints, silently, forever.
**Fix #10448** ("[Fix][Zeta] make the job failed when triggering checkpoint
fails (apache#10442)")
replaces both catch blocks in that method with a call to
`handleCoordinatorError(...,
CheckpointCloseReason.CHECKPOINT_INSIDE_ERROR)`.
### Current `dev` behavior this test asserts (verified against the current
HEAD, not assumed)
I traced the full current call chain in `CheckpointCoordinator.java`,
`CheckpointManager.java`, `JobMaster.java`, `SubPlan.java`, and
`PhysicalPlan.java` before writing
this test:
- `startTriggerPendingCheckpoint`'s `catch (Exception e)` (around
`CheckpointCoordinator.java:965-970`) calls
`handleCoordinatorError("triggering checkpoint barrier failed", e,
CheckpointCloseReason.CHECKPOINT_INSIDE_ERROR)`.
- `handleCoordinatorError` (`CheckpointCoordinator.java:341-357`) marks the
coordinator `FAILED`,
calls `checkpointManager.handleCheckpointError(pipelineId, false)`, and
resets `pendingCounter`
to 0 via `cleanPendingCheckpoint`.
- `CheckpointManager#handleCheckpointError` ->
`JobMaster#handleCheckpointError`
(`JobMaster.java:767-779`) -> `SubPlan#handleCheckpointError()`
(`SubPlan.java:639-646`), which
cancels the pipeline (`PipelineStatus.CANCELING`) if it is not already in
an end state.
- Once all tasks finish canceling, `SubPlan#getPipelineEndState()`
(`SubPlan.java:243-290`) sees
`canceledTaskNum > 0` and calls `cancelCheckpoint()`; because the
coordinator's own
`checkpointCoordinatorFuture` was already completed `FAILED` by
`handleCoordinatorError`, that
call returns the same `FAILED` status, which upgrades the pipeline's end
state from `CANCELED`
to `FAILED`.
- With `job.retry.times = 0` (this test's conf),
`SubPlan#canRestorePipeline()` is `false`
(`getPipelineRestoreNum() < pipelineMaxRestoreNum` is `0 < 0`), so
`stateProcess()`'s
`FAILED`/`CANCELED` case skips restore and completes the pipeline future
as `FAILED`.
- `PhysicalPlan#addPipelineEndCallback` (`PhysicalPlan.java:148-196`) sees
the (single) pipeline
end as `FAILED`, so `failedPipelineNum > 0` and the whole job transitions
to
`JobStatus.FAILED` via `updateJobState`.
So the current, documented behavior is: **the job reaches a terminal
`FAILED` state**, not the
old silent-forever-`RUNNING` hang. This is exactly what the merged fix's own
unit test
(`CheckpointBarrierTriggerErrorTest`) already asserts with mocks; this PR
adds equivalent coverage
at the E2E level with a real, unmocked, unmodified `queryTaskGroupAddress`.
### Trigger mechanism, and why it is reliable without having run it locally
Per this project's Apache SeaTunnel local-verification policy, I have not
run this test locally;
GitHub CI on this PR's head is the first execution. I was correspondingly
rigorous about tracing
every step of the mechanism against the real, current source before writing
the assertions:
`CheckpointCoordinator#triggerCheckpoint`
(`CheckpointCoordinator.java:1120-1132`) is the only code
that can make `startTriggerPendingCheckpoint`'s
`CompletableFuture.allOf(completableFutureArray)
.get()` throw *synchronously* (as opposed to a per-task RPC merely failing
later, which that
`allOf` call never even notices, since it only waits for
`triggerCheckpoint()` to return -- a
separate, already-covered dead-letter scenario, see #12111). It maps every
starting subtask through
`checkpointManager::sendOperationToMemberNode`
(`CheckpointManager.java:386-400`), which calls
`jobMaster.queryTaskGroupAddress(...)` (`JobMaster.java:977-994`) *before*
issuing the RPC. That
method does exactly one thing that can throw:
`ownedSlotProfilesIMap.get(pipelineLocation)`
returning `null`, which throws `IllegalArgumentException("can't find task
group address from
taskGroupLocation: ...")`.
I confirmed (via a repo-wide search) that `ownedSlotProfilesIMap`'s only
entry-removal call site
(`JobMaster#releasePipelineResource`, `JobMaster.java:923-950`) runs only
after a pipeline has
already left `RUNNING`, by which point `cleanPendingCheckpoint` has already
cancelled that
coordinator's own scheduler (`scheduler.shutdownNow()`,
`CheckpointCoordinator.java:1203`) -- so
nothing in the running system naturally races this lookup against a live,
scheduled trigger.
Killing/isolating a worker (this test class's usual technique) does not
reach this code path
either: a graceful leave fails the *task* directly via
`CoordinatorService#failedTaskOnMemberRemoved` without ever touching this
map, while an ungraceful
one leaves a *stale but present* entry (the RPC fails later, asynchronously
-- the separate
dead-letter case covered by #12111, not this one).
So this test uses a different, still entirely real, lever instead of cluster
membership:
`ownedSlotProfilesIMap` is a plain, named Hazelcast `IMap`
(`Constant#IMAP_OWNED_SLOT_PROFILES`),
obtained the exact same way this test class's own `getReadyToCloseCount`
helper already reads
`Constant#IMAP_RUNNING_JOB_STATE` directly, and the same way the
engine-server module's own
`EngineStateStoreMetricExportsTest` pokes this exact map in its unit tests.
The test removes this
job's entry from that live, shared map -- this is not a mock and not a
reflected exception injected
into production code: it is the same real, unmodified, running
`queryTaskGroupAddress` that throws
its own real `IllegalArgumentException` the next time it executes, exactly
as it would if this
bookkeeping ever went missing for any other reason. I checked every other
reader of this map
(metrics export, pipeline cleanup,
`PhysicalVertex#checkTaskGroupIsExecuting` -- itself only
reachable via master-failover restore, never steady-state `RUNNING`) and
confirmed they all
null-check and skip gracefully, so the removal cannot trip any other code
path first.
This is deterministic, not a narrow-window race like a worker kill: the
entry is left removed
permanently (the pipeline is about to fail anyway), so the very next
scheduled trigger attempt
that has not already started is guaranteed to observe the missing entry once
the removal
completes -- no timing window to miss.
### Test outline
`CheckpointCoordinatorFailoverIT#testStreamJobFailsAfterCheckpointTriggerDispatchFailure`:
1. Starts a single embedded node running a `parallelism = 1` streaming
FakeSource -> LocalFile job
(`checkpoint.interval = 2000`, `job.retry.times = 0`).
2. Waits for the job to be `RUNNING` and producing rows.
3. Waits for the checkpoint-id counter to reach `2`, which (since
`tryTriggerPendingCheckpoint`
never allocates a new id while `pendingCounter > 0`) can only happen once
checkpoint id 1 has
fully completed -- proving checkpointing was healthy before the fault is
injected.
4. Removes the job's entry from `ownedSlotProfilesIMap` (the real-fault
injection described above).
5. Asserts the job reaches terminal `JobStatus.FAILED` within a bounded
window, and that
`JobResult#getError()` contains
`CheckpointCloseReason.CHECKPOINT_INSIDE_ERROR`'s message text
(traced end-to-end through `SubPlan`/`PhysicalPlan` back to the
coordinator's own error).
The bounded `Awaitility` wait for `FAILED` also implicitly demonstrates the
pre-fix bug no longer
occurs: against the old code, the job would stay `RUNNING` forever with no
further checkpoints, so
this wait would time out and fail the test.
### Does this PR introduce any user-facing change?
No. Test-only; no `src/main` changes.
### How was this patch tested?
- `./mvnw spotless:apply -pl
seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base -nsu
-Dmaven.gitcommitid.skip=true` -- BUILD SUCCESS, no further formatting changes
needed.
- Per this repo's Apache SeaTunnel local-verification policy, no local
compile/test/E2E execution
was performed; GitHub CI on this PR's head is the authoritative
verification.
### Check list
- [x] Test-only change, no `src/main` modifications.
- [x] No new dependencies.
- [x] No documentation changes needed (test-only).
🤖 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]