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]

Reply via email to