DanielLeens commented on PR #12152:
URL: https://github.com/apache/seatunnel/pull/12152#issuecomment-5601156434
Thanks @SEZ9 — first, to close the "cut off" question: I checked my own
comment's raw body and it isn't truncated, it ends on a complete sentence
("...not a specific implementation shape."). I think this is the same GitHub
"show more" collapse on long comments that tripped us both up on #12130 the
other day, not an actual truncation. For the record, my proposed shape there
was: log at WARN with the master-switch/shutdown context, then complete
`checkpointCoordinatorFuture` and hand off to
`checkpointManager.handleCheckpointError(...)`, skipping anything that assumes
it's off the operation thread — but see below, I don't think that shape is
right either once you dig one level further.
On your actual question — whether a visible operation failure is an
acceptable fallback for the shutdown-triggered rejection case — I went and
checked rather than guessing, and I don't think it is safe, for a reason
neither of us had raised yet:
**Letting the exception escape does not actually get anyone out of
`.join()`.**
- `clearCoordinatorService()` (`CoordinatorService.java:1281-1312`) calls
`executorService.shutdownNow()` (`:1312`), which sends `Thread.interrupt()` to
whatever is running on that pool.
- But the consumer side of `checkpointCoordinatorFuture` blocks via
`.join()` — e.g. `SubPlan.getPipelineEndState()` at `SubPlan.java:271-275`,
`jobMaster.getCheckpointManager().waitCheckpointCoordinatorComplete(getPipelineId()).join()`.
`CheckpointManager.waitCheckpointCoordinatorComplete`
(`CheckpointManager.java:420-423`) returns a `PassiveCompletableFuture` chained
onto `checkpointCoordinatorFuture`, and both `PassiveCompletableFuture` and the
project's own
`org.apache.seatunnel.engine.common.utils.concurrent.CompletableFuture`
(`seatunnel-engine-common/.../concurrent/CompletableFuture.java`) implement
`join()` as a bare `return super.join()` — plain JDK `CompletableFuture.join()`
semantics.
- JDK `CompletableFuture.join()` does not honor `Thread.interrupt()`. Unlike
`get()`, its internal `Signaller` is built with `interruptible=false`, so an
interrupted waiter just keeps parking until the future actually completes.
- So a thread parked in that `.join()` is not released by `shutdownNow()`'s
interrupt, and nothing else completes `checkpointCoordinatorFuture` for it on
the master-step-down path — `JobMaster.interrupt()`
(`JobMaster.java:1502-1505`) only flips `isRunning` and completes the unrelated
`jobMasterCompleteFuture`, it doesn't touch `checkpointCoordinatorFuture`.
Net effect: in the shutdown-triggered rejection window, "let the operation
fail visibly and rely on teardown" leaves that waiter parked indefinitely
anyway — the same hang this PR is trying to remove, just relocated. It also
means `clearCoordinatorService()`'s own `executorService.awaitTermination(20,
TimeUnit.SECONDS)` (`CoordinatorService.java:1316`) would likely time out
(logged as a warning, not fatal, but a real leaked thread). This isn't new to
this PR — it's a pre-existing property of how `.join()` is used against this
future — but it does mean the two rejection causes don't actually need
different treatment: saturation and shutdown both leave a waiter that only the
future's completion can free.
Given that, I'd narrow the fallback further than what either of us proposed
before: on `RejectedExecutionException`, don't call the full
`handleCoordinatorError(...)` chain (that's still the
`checkpointManager.handleCheckpointError(pipelineId, false)` →
`SubPlan.handleCheckpointError()` → `synchronized
updatePipelineState(CANCELING)` → `stateProcess()`/`task.cancel()` traversal we
already agree can't run on the operation thread). Instead, just complete
`checkpointCoordinatorFuture` directly to a FAILED state in the catch block.
That's a private field on this class so `reportCheckpointErrorFromTask` already
has direct access. Completing it is O(1) — it flips the future's internal state
and unparks any `.join()`/`.get()` waiters via `LockSupport.unpark`, it does
not run the waiter's downstream code on the completing thread, and I didn't
find any `.thenApply`/`.whenComplete` continuation attached to
`checkpointCoordinatorFuture` itself inside `CheckpointCoordinator` (only out
side consumers wrapping it in a `PassiveCompletableFuture` and `.join()`-ing),
so there's no hidden reentrant chain riding along on this thread either.
The tradeoff I can't fully resolve from here: this unblocks the waiter with
the correct terminal status, but it skips the active-cancellation side effects
that `checkpointManager.handleCheckpointError(...)` normally drives (cancelling
still-running tasks, etc.). Whether that's acceptable probably depends on
whether, by the time this rejection can actually happen (pool saturated on a
live master, or torn down on a step-down), those tasks are already being
handled through some other path — that's the part I'd want @CryoThrust's read
on, since it touches restore/cancellation sequencing I haven't traced end to
end for this PR. If tasks genuinely need active cancellation here too, then the
future-completion alone isn't sufficient and we're back to needing an async
best-effort retry (e.g., a short bounded resubmission attempt, or scheduling
onto a small dedicated single-thread "last resort" executor that outlives the
coordinator pool) rather than anything synchronous on the operation
thread.
So, concretely, what I'd like to see: complete `checkpointCoordinatorFuture`
inline on `RejectedExecutionException` (cheap, unblocks joiners, no
operation-thread state-machine work), plus a WARN log carrying
jobId/pipelineId/reason, and then a test that pins exactly this — a coordinator
built with a pre-shutdown executor, asserting the future still resolves and any
`.join()` caller unblocks, rather than only asserting the exception doesn't
propagate.
--
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]