Vivek1106-04 commented on issue #12339: URL: https://github.com/apache/seatunnel/issues/12339#issuecomment-5750443171
Thanks for looking at this closely. I want to push back a bit, because I think the code puts the scaling concern in a different pool than the one the comment targets. **Where the network work actually runs** The shared scheduler deliberately does not carry the RPC path: - **Timer threads** only re-arm the trigger / timeout watchdog and hand the task to the dispatch pool. They never touch the network. - **Dispatch threads** run exactly two bodies. `tryTriggerPendingCheckpoint` takes the per-coordinator `lock`, does state checks, and calls `createPendingCheckpoint` -> `triggerPendingCheckpoint`, which is `supplyAsync(..., executorService).thenApplyAsync(..., executorService)`. The timeout watchdog is a `pendingCheckpoints` lookup plus an error path. Both return quickly. - **The barrier trigger and ACK handling** run on `executorService`, the coordinator-service pool, via `thenApplyAsync(this::triggerCheckpoint, executorService)`. So when workers, jobs, and pipelines grow, the pressure lands on the coordinator-service pool, not on the shared timer/dispatch pools. **That pool is already configurable** `engine.coordinator-service.core-thread-num` (default 10) and `max-thread-num` (default `Integer.MAX_VALUE`) over a `SynchronousQueue`: it already grows without bound, and it is already the documented dial for exactly this scaling story. Adding `checkpoint.scheduler-*-thread-num` gives operators a second dial for one symptom, and the newer, more checkpoint-sounding name is the one they will reach for first -- while the pool that is actually saturated is the other one. **What the shared pools do today** Dispatch is `max(8, availableProcessors * 2)`, threads created on demand, reaped after 60s idle, unbounded queue that never rejects (a rejection here would be a dropped trigger or a dropped watchdog). Undersizing it yields delay, not loss, and that delay is only reachable if all dispatch threads block at once -- which requires the dispatched bodies to block, and they do not. **The one real hole, which a knob would paper over** `startTriggerPendingCheckpoint` registers `pendingCompletableFuture.thenAccept(...)`, and that body contains a blocking `CompletableFuture.allOf(completableFutureArray).get()`. Normally the pending future is still incomplete (its chain runs on `executorService`), so the body lands on a coordinator-service thread. But if it has already completed, `thenAccept` runs inline on the dispatch thread and blocks it on the barrier send. The fix for that is `thenAcceptAsync(..., executorService)`, not a bigger dispatch pool. I would rather close that hole than add a dial that hides it. **What I propose instead -- invert the order of your two comments** Metrics first, option second, and only if the metrics show it: 1. Ship this PR with the current defaults. 2. Follow-up: the metrics from your second comment -- dispatch queue depth, active dispatch threads, and **scheduling delay** (actual fire time minus scheduled fire time). Scheduling delay is the discriminator: sustained nonzero means the pool is undersized; flat zero means the pool is irrelevant to whatever the operator is chasing. 3. If a real separated deployment then shows sustained delay, add the option with documentation naming the metric to tune it against. Option names are permanent contracts here. I would rather not ship one for a hypothetical, and today an operator has no signal to tune it against -- which is precisely your second comment. A dial without a gauge is guessing, and guessing low silently delays triggers and watchdogs for every pipeline on the member. **If you still want it in this PR** I will add it, but as one option, not two: `checkpoint.scheduler-dispatch-thread-num`, default `0` meaning auto (`max(8, availableProcessors * 2)`), so no existing deployment changes behavior. No timer option -- timer threads only hand off, and there is no signal an operator could tune that against. And a direct question: do you have a deployment where checkpoint triggers already drift under load? If there is a reproduction, that settles it and I will size from the measurement rather than from the argument above. -- 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]
