Vivek1106-04 commented on issue #12339:
URL: https://github.com/apache/seatunnel/issues/12339#issuecomment-5757246015

   Agreed on both.
   
   **1. Dropping `checkpoint.scheduler-dispatch-thread-num`.**
   
   Not adding it in this PR. The dispatch runnables do no I/O and no RPC — they
   allocate a checkpoint ID, build the `PendingCheckpoint` and hand off to the
   coordinator's own executor — so with the blocking `allOf(...).get()` moved 
off
   the dispatch thread by `thenAcceptAsync`, there is no longer a case I can 
point
   to where the auto-sized pool is wrong. If the scheduling-delay metric later
   shows dispatch queueing under real load, the option comes back with evidence
   behind it.
   
   **2. Benchmark — agreed, separate PR. One design note so it measures the 
right thing.**
   
   You are right that `checkpointSingleInput` is blind to this. 
`CheckpointBenchmarkTrigger`
   reflects straight into `createPendingCheckpoint` + 
`startTriggerPendingCheckpoint`,
   so both the old `ScheduledExecutorService` and the new 
`PipelineCheckpointScheduler`
   lease are out of the measured path entirely.
   
   The thing to be careful about is that the obvious extension — route the 
bridge
   through the scheduler and re-measure `checkpointSingleInput` — would still 
not
   separate the implementations. That number is dominated by the barrier ACK
   round-trip; the scheduler contributes one queue hop on top of it. At one
   pipeline on an idle host, old and new will look the same, and that result 
would
   be read as "no difference" when the benchmark simply cannot see the axis the
   change is on.
   
   The axis is **pipeline count**. On dev, each `CheckpointCoordinator` builds 
its
   own pool:
   
   ```java
   this.scheduler = Executors.newScheduledThreadPool(2, ...);   // per (jobId, 
pipelineId)
   ```
   
   so a member running P pipelines carries 2P timer threads. The shared 
scheduler
   is `TIMER_THREAD_NUM` timer threads plus a bounded dispatch pool
   (`max(8, cores * 2)`), independent of P. So what a benchmark should report, 
as
   P sweeps 1 / 10 / 100 / 500:
   
   - **scheduling delay** — scheduled fire time vs actual fire time, per 
pipeline,
     as a distribution (p50/p99/max), not a mean. This is the 
correctness-adjacent
     number: if the shared pool ever delays a trigger past `checkpointMinPause` 
or
     the timeout window, that shows up here and nowhere else.
   - **thread count and thread CPU** for the scheduling machinery, which is 
where
     the shared model is expected to win and the per-coordinator model to 
degrade
     linearly.
   - **dispatch queue depth** at steady state, which is the same signal the
     follow-up metrics issue would expose at runtime — worth having both come 
from
     one definition so the benchmark and production metric agree.
   
   That shape is a latency-under-concurrency measurement, so `Mode.SampleTime` 
on
   the fire-delay rather than `Mode.AverageTime` on completion, with the 
pipelines
   as `@State` setup rather than the benchmarked call. I would keep
   `checkpointSingleInput` exactly as it is alongside it — unchanged, so the
   single-pipeline completion path stays comparable across the change — and add 
the
   scheduler benchmark as a separate class rather than extending that one.
   
   On sequencing: I will land the benchmark first, as its own PR against dev, so
   the numbers exist before #12165 is judged on them. Running it on dev gives 
the
   per-coordinator baseline; running the same benchmark on the #12165 branch 
gives
   the comparison. If it shows the shared scheduler is worse on scheduling 
delay at
   any P, that is a result worth having before the merge decision, not after.
   


-- 
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