DanielLeens commented on PR #12152:
URL: https://github.com/apache/seatunnel/pull/12152#issuecomment-5578517291
Thanks @SEZ9 — I traced both points against the current head (`f3fcb611e4`)
and can confirm your read on both, plus have a concrete answer to the
parenthetical question.
**Point 1 (rejection fallback) — confirmed still open.**
`reportCheckpointErrorFromTask()` (`CheckpointCoordinator.java:518-528`) is
unchanged from what I originally reviewed: a bare
`executorService.execute(...)` with nothing catching
`RejectedExecutionException`. Agreed this is still a live gap.
**Point 2 (test coverage) — confirmed still a gap.**
`testCheckpointErrorReportDoesNotRunOnCallerThread`
(`CheckpointCoordinatorTest.java:979-1007`) only proves the call returns
immediately while the single-thread executor is busy
(`Mockito.verifyNoInteractions(checkpointManager)` right after the call, before
releasing the latch). There's no assertion that the error is actually delivered
once the executor frees up, and no case at all with an
already-shutdown/rejecting executor. Agreed both need to be added.
**On your parenthetical ("re-entering the JobMaster state machine on the
operation thread during shutdown") — this is a real concern, not a
hypothetical.** I checked the caller:
`CheckpointErrorReportOperation.runInternal()`
(`CheckpointErrorReportOperation.java:42-49`) extends `TaskOperation`, a
Hazelcast `Operation`, and calls `reportCheckpointErrorFromTask()`
synchronously from `runInternal()`. So the method's only caller genuinely is a
Hazelcast operation thread. Because `ExecutorService.execute()` throws
`RejectedExecutionException` synchronously to its caller (it doesn't route
through any future), a naive `try { executorService.execute(...) } catch
(RejectedExecutionException e) { handleCoordinatorError(...); }` fallback would
run `handleCoordinatorError()` — which reaches into
`checkpointManager.handleCheckpointError(...)` and
`cleanPendingCheckpoint(...)` — directly on that same operation thread. That's
exactly the class of work this PR's title is trying to move off
the operation thread, just resurfacing in the one narrow window (post-shutdown
/ pool-saturated) where the async handoff itself fails.
One more thing worth flagging while we're here: mirroring `reportedTask()`'s
`CompletableFuture.runAsync(..., executorService).exceptionally(...)` shape
(`CheckpointCoordinator.java:311-332`) will **not** actually catch a rejected
submission either. `CompletableFuture.runAsync`/`supplyAsync` call
`executor.execute(...)` directly inside the static factory method, before the
`CompletableFuture` is even returned — if that `execute()` throws, the
exception propagates synchronously out of `runAsync()` to its caller, it is
never routed into the returned future's completion, so `.exceptionally(...)`
never sees it. That means `reportedTask()` itself likely has the same latent
hole against a torn-down executor today (separate issue, out of scope for this
PR, but worth a follow-up ticket). The practical implication for this PR:
whichever fallback shape gets picked, the `RejectedExecutionException` guard
has to wrap the raw `execute()`/`runAsync()` call directly in a `try/catch` —
an `.e
xceptionally()` continuation alone will not do it.
Given that, my suggestion: keep the fallback but make it deliberately
minimal rather than a full inline traversal — e.g. log at WARN with the
master-switch/shutdown context, then complete `checkpointCoordinatorFuture` and
hand off to `checkpointManager.handleCheckpointError(...)` the same way
`handleCoordinatorError()` does, but skip re-entering anything that assumes
it's off the operation thread. `handleCoordinatorError()` already early-returns
via `checkpointCoordinatorFuture.isDone()`, so re-entrancy during the shutdown
race isn't the issue — the issue is purely "is it acceptable to do this bounded
amount of work on the operation thread in the rare rejection case." I'd say
yes: a short, documented, best-effort inline completion on rejection is
strictly better than silently dropping the report and leaving the pipeline
parked on `.join()` forever, as long as it's called out in a comment so it
isn't mistaken for the common path. Happy to see it land either way — the
blocking
requirement from my side is just that the guard exists and is tested
(positive-delivery-after-release case, plus a shutdown/rejected-at-call-time
case), not a specific implementation shape.
--
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]