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]