davidzollo opened a new pull request, #12111:
URL: https://github.com/apache/seatunnel/pull/12111
## What this PR verifies and covers
`CheckpointCoordinator#startTriggerPendingCheckpoint` triggers a checkpoint
barrier like this
(`seatunnel-engine/seatunnel-engine-server/.../checkpoint/CheckpointCoordinator.java:944-958`):
```java
CompletableFuture<InvocationFuture<?>[]> completableFutureArray =
CompletableFuture.supplyAsync(() -> new CheckpointBarrier(...),
executorService)
.thenApplyAsync(this::triggerCheckpoint, executorService);
CompletableFuture.allOf(completableFutureArray).get();
```
`completableFutureArray` is a single `CompletableFuture` whose *value* is an
`InvocationFuture<?>[]` -- not the array itself.
`CompletableFuture.allOf(...)` only spreads a
real array argument into its varargs; handed one future, it waits on exactly
that future, i.e.
until `triggerCheckpoint()` returns (the per-task
`CheckpointBarrierTriggerOperation` RPCs are
*fired*), not until any of them land or are acknowledged. The individual
`InvocationFuture`s
inside the resolved array are never consulted again by anything in this
class, so a barrier-RPC
failure to a since-unreachable worker is silently discarded -- a true dead
letter, not merely
delayed handling.
Contrast the same class's own correct pattern for `notifyTaskStart()` /
`notifyCompleted()`
(lines ~451-452, ~474-481): both spread an already-resolved
`InvocationFuture<?>[]` directly into
`allOf`, which genuinely waits on every element.
I independently re-verified this against the current `dev` HEAD before
writing the test --
confirmed the exact code, line numbers, and contrasting patterns cited above
are still present
and unchanged.
The only backstop for a barrier-dispatch RPC that never lands is the
scheduled per-pending-checkpoint
timeout later in the same method (lines ~972-1001): once
`checkpoint.timeout` elapses without the
checkpoint becoming fully acknowledged,
`CheckpointCloseReason#CHECKPOINT_EXPIRED` fires, cancelling
and restarting the pipeline the same way a hard task failure would.
## Test added
`CheckpointCoordinatorFailoverIT#testStreamJobRecoversAfterWorkerUnreachableDuringCheckpointBarrierDispatch`
- Starts a 1-master/2-worker cluster and runs a parallelism=1 streaming job
(single pipeline,
single task) so the target worker can be identified unambiguously.
- Identifies which worker hosts the task via its live execution address,
then tightly polls
(10ms) the shared checkpoint-id counter state store used elsewhere in this
file, terminating the
target worker in the same loop iteration that first observes a new
checkpoint id -- i.e. as soon
as that checkpoint's barrier dispatch is imminent or already starting.
Termination uses
`HazelcastInstance.getLifecycleService().terminate()` (not `shutdown()`),
which per Hazelcast's
own `Node#shutdown(boolean terminate)` skips the graceful cluster-leave
notice that would
otherwise let `CoordinatorService#failedTaskOnMemberRemoved` fail the task
almost immediately --
that fast, generic path is exactly what the existing
`ClusterFaultToleranceIT`-style worker-kill
tests exercise, and it would mask this specific dispatch-window bug
entirely.
- The job's `checkpoint.timeout` (8s) is deliberately far below Hazelcast's
own native membership
failure-detection ceiling. That ceiling is *not* incidental here: this
build's shaded Hazelcast
jar compiles in a 60-second default for
`hazelcast.max.no.heartbeat.seconds` (verified by
disassembling `com.hazelcast.spi.properties.ClusterProperty`'s static
initializer), and this
test module's own `hazelcast.yaml` does not raise it -- only this repo's
top-level,
production-only `config/hazelcast.yaml` does, and that file is not on this
module's test
classpath. Rather than depend on whichever `hazelcast.yaml` happens to be
on the classpath, the
test explicitly overrides `hazelcast.max.no.heartbeat.seconds` (to 180s)
on every node it starts,
so Hazelcast's own native failure detection provably cannot fire inside
the test's bounded
recovery wait, regardless of ambient config now or in the future.
- Asserts the job recovers: status back to RUNNING, the task redeployed onto
the surviving worker
(by execution address), and row output resuming -- all within a bounded
window tied to
`checkpoint.timeout`, not Hazelcast's much slower native detection.
**What this test proves:** the job recovers within a bounded time tied to
`checkpoint.timeout`,
even though the coordinator cannot detect this specific barrier-dispatch RPC
failure immediately.
**What it does NOT prove:** that the barrier-dispatch RPC failure is caught
immediately -- per the
bug above, it is not, and this test does not assert instant detection. It
also does not modify
`CheckpointCoordinator` itself; this is test-only coverage of
currently-fragile behavior.
## Test plan
- `./mvnw spotless:apply -pl
seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base -am -nsu` --
passed.
- Compilation verified locally: `./mvnw install -pl
'seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base' -nsu
-Dmaven.gitcommitid.skip=true -DskipTests -Dspotless.check.skip=true` -- BUILD
SUCCESS, with the compiled `CheckpointCoordinatorFailoverIT.class` (containing
the new test method) confirmed present under `target/test-classes`. This
narrower, non-`-am` invocation resolved all upstream dependencies from a local
repository already freshly built from this same unmodified `dev` base (this
machine was running many concurrent, unrelated builds during verification); a
full `-am` invocation from a clean state is expected to succeed identically
since no `src/main/**` code was touched.
- Actual test execution: per this repo's Apache SeaTunnel local-verification
policy, this PR relies on GitHub CI (the fork's Actions run behind the `Build`
check) as the authoritative execution result; local test execution was not
performed.
🤖 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]