davidzollo opened a new pull request, #12314:
URL: https://github.com/apache/seatunnel/pull/12314
## Summary
`CheckpointCoordinator#startTriggerPendingCheckpoint` fires every per-task
checkpoint barrier RPC via `triggerCheckpoint()` (which returns
`InvocationFuture<?>[]`, `CheckpointCoordinator.java:1167-1179`), but only ever
waited for:
```java
CompletableFuture.allOf(completableFutureArray).get();
```
`completableFutureArray` is `CompletableFuture<InvocationFuture<?>[]>` - the
future that resolves once `triggerCheckpoint()` has *produced* the array, i.e.
once every barrier RPC has been *sent* via `sendOperationToMemberNode`. Because
`InvocationFuture` extends the JDK `java.util.concurrent.CompletableFuture`
(`InvocationFuture -> AbstractInvocationFuture -> InternalCompletableFuture ->
java.util.concurrent.CompletableFuture`), the type system silently accepted
wrapping that single array-producing future as one varargs element to
`allOf(...)`. The per-task RPCs themselves - the actual `InvocationFuture`
elements inside the array - were never observed anywhere else in the method
(the array is used only to build this one `allOf(...)` call).
Consequence: a barrier RPC that fails permanently (target member left, or a
non-retryable Hazelcast operation failure) is a dead letter. The only backstops
are the `checkpoint.timeout` scheduled task
(`CheckpointCoordinator.java:1019-1047` after this change, previously
~`:972-1001`) which reports a generic `CHECKPOINT_EXPIRED` only after the full
timeout elapses, and member-removal handling. This is bounded (not P0) but
delays detection and hides the real failure cause behind a generic timeout.
This was raised during review of #12111 (E2E coverage for this exact
dead-letter scenario); the reviewer requested a production fix rather than
test-only coverage.
## Fix
In `startTriggerPendingCheckpoint` (the only method touched,
`CheckpointCoordinator.java`):
1. Keep the existing `.get()` on `completableFutureArray` unchanged in
behavior (same catch blocks, same `CHECKPOINT_INSIDE_ERROR` handling for the
synchronous-throw path used by #10448) - now capturing the returned array into
`InvocationFuture<?>[] barrierFutures`.
2. Observe the real per-task futures **asynchronously and non-blocking**:
`CompletableFuture.allOf(barrierFutures).whenCompleteAsync((ignored, throwable)
-> {...}, executorService)`.
- Non-blocking is essential: the `checkpoint.timeout` scheduled task
right after this block is the only backstop for an RPC that never completes at
all (member hang, network partition before Hazelcast's own retries give up); a
blocking join here would prevent that timeout from ever getting scheduled.
- `executorService` is ultimately the `SynchronousQueue`-backed pool
created by `CoordinatorService#createCoordinatorExecutor`
(`CoordinatorService.java:287-298`), which rejects new work once its threads
are busy - a blocking wait on this thread could starve other pipelines'
checkpoint coordination sharing the same pool.
3. On a real failure, guard against a stale callback: by the time a slow RPC
finally fails, the checkpoint may already have completed and been removed from
`pendingCheckpoints` (`completePendingCheckpoint`,
`CheckpointCoordinator.java:1420`), or the whole coordinator may already be
closed. The callback checks `pendingCheckpoints.containsKey(checkpointId)`
(mirroring the same idiom already used by the `checkpoint.timeout` task) before
calling `handleCoordinatorError`, so a late failure belonging to an
already-completed checkpoint is a no-op.
4. Route a real failure to `handleCoordinatorError(...,
CheckpointCloseReason.CHECKPOINT_INSIDE_ERROR)` - the same reason already used
for the synchronous-throw path in this method (the #10448 path), so the
coordinator now reports the actual dispatch failure immediately instead of a
generic `CHECKPOINT_EXPIRED` after the full timeout.
Blast radius: every checkpoint and savepoint trigger goes through this
method, so this changes the failure-detection latency/cause for *any* barrier
dispatch RPC failure across all pipelines - previously silent-until-timeout,
now immediate with the real cause. No behavior changes when barrier RPCs
succeed.
## Test-double fix (required by this change)
`CheckpointCoordinatorTest`'s `TestCheckpointManager` and a couple of inline
`Mockito.mock(CheckpointManager.class)` stubs (`buildMinimalCoordinator`)
return `null` from `sendOperationToMemberNode(...)`. That `null` previously
never reached `allOf(...)` since only the array-producing future was observed;
with this fix, `CompletableFuture.allOf(barrierFutures)` would
`NullPointerException` synchronously on a `null` array element, and since that
call sits inside a `thenAccept(...)` continuation whose result is never read,
the exception would be silently swallowed - hanging the checkpoint forever with
no error surfaced.
I audited every existing test that reaches `startTriggerPendingCheckpoint`
for real (not via a `Mockito.doNothing()`-stubbed
`tryTriggerPendingCheckpoint`):
- `testFilteringClosedTasksAndActions` uses `TestCheckpointManager` (returns
`null`), but calls `coordinator.triggerCheckpoint(barrier)` **directly** - it
never goes through `startTriggerPendingCheckpoint`, so it's unaffected.
- `testSchedulerThreadShouldNotBeInterruptedBeforeJobMasterCleaned`,
`testCheckpointContinuesWorkAfterClockDrift`, `testCheckpointMinPause` all use
a real (or real-delegating) `CheckpointManager` whose
`sendOperationToMemberNode` dispatches a genuine (non-null) Hazelcast
`InvocationFuture`.
- No existing test's happy path actually feeds a `null` element into the new
`allOf(barrierFutures)` call. No test-double changes were needed to keep
existing tests passing.
## New tests (`CheckpointCoordinatorTest.java`)
- **`testBarrierDispatchFailureIsRoutedToCoordinatorError`**: stubs
`sendOperationToMemberNode` to return an already `completeExceptionally`'d
`InvocationFuture`. Since `InvocationFuture` is `final` with a package-private
constructor (cannot be subclassed or instantiated from test code), it is mocked
via Mockito's inline mock maker (already enabled for this module via
`src/test/resources/mockito-extensions/org.mockito.plugins.MockMaker` =
`mock-maker-inline`, and already used elsewhere in this file for the same class
and for `Mockito.mockStatic`). `completeExceptionally` is stubbed with
`doCallRealMethod()` so it genuinely runs the inherited JDK
`CompletableFuture#completeExceptionally` on that instance, setting its real
internal completion state - which `CompletableFuture.allOf(...)`'s own
internals read directly. The checkpoint timeout is configured to 10 minutes,
far longer than the test's 10-second Awaitility bound, so a pass can only be
explained by the new asynchronous path,
not the timeout backstop. Asserts `handleCoordinatorError` is invoked with
`CHECKPOINT_INSIDE_ERROR` and the real cause (allowing for the
`CompletionException` wrapping that `CompletableFuture.allOf(...)` applies per
its documented contract).
- **`testStaleBarrierDispatchFailureForCompletedCheckpointIsIgnored`**: same
failing RPC, but the checkpoint id is deliberately never registered in
`pendingCheckpoints` (simulating a checkpoint that already completed and was
removed). Asserts `handleCoordinatorError` is never called within a bounded
wait - the stale-callback guard.
## Verification
- `./mvnw spotless:apply -pl seatunnel-engine/seatunnel-engine-server -nsu
-Dmaven.gitcommitid.skip=true` - clean, no additional formatting changes.
- Per this task's constraints, no local compile, test, or other build/run
verification was performed. Every symbol, signature, and generic type used in
the production change and the new tests was manually verified against the
pinned source at this PR's base commit (including a `javap` inspection of the
shaded Hazelcast jar to confirm `InvocationFuture`'s exact class hierarchy and
that it is `final`). GitHub CI on this PR's head is the verification of record
for compilation and test execution.
## Test plan
- [ ] GitHub CI (compile + `seatunnel-engine-server` unit tests) passes on
this PR's head.
🤖 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]