Vivek1106-04 opened a new issue, #12339:
URL: https://github.com/apache/seatunnel/issues/12339

   # Member-wide shared checkpoint trigger scheduling for Zeta
   
   Improvement issue: 
[#12122](https://github.com/apache/seatunnel/issues/12122).
   Candidate implementation: [PR 
#12165](https://github.com/apache/seatunnel/pull/12165).
   
   This proposal is raised at the request of @nzw921rx on PR #12165, so that 
the scheduling model is agreed on the dev list before the implementation is 
considered for merge. Nothing in the linked PR should be read as an accepted 
design.
   
   ## Background
   
   `CheckpointCoordinator` constructs its own 
`Executors.newScheduledThreadPool(2, ...)` in its constructor, and 
`CheckpointManager` constructs one coordinator per entry of the checkpoint plan 
map, that is, one per pipeline. Timer threads therefore scale linearly with the 
number of checkpoint-enabled pipelines on a member. A member running 500 such 
pipelines holds up to 1000 threads that are asleep nearly all of the time.
   
   The pool is only ever used as a timer. There are exactly two scheduling call 
sites in `CheckpointCoordinator`:
   
   1. re-arming the periodic trigger, in `scheduleTriggerPendingCheckpoint`, 
which runs `tryTriggerPendingCheckpoint`
   2. arming the per-checkpoint timeout watchdog, inside 
`startTriggerPendingCheckpoint`
   
   The heavy barrier work already runs on the injected `executorService` via 
`thenApplyAsync`, not on the scheduled pool. That is the property that makes a 
shared timer viable at all.
   
   `cleanPendingCheckpoint` currently calls `shutdownNow()` on the 
per-coordinator pool and constructs a replacement, which means the pool's 
lifecycle is bound to a single pipeline's cleanup.
   
   Two details are worth stating because they bound the size of the actual 
problem:
   
   - Core workers are created on demand; there is no `prestartAllCoreThreads` 
anywhere in `seatunnel-engine`. Periodic triggering begins from `allTaskReady`, 
and `restoreCoordinator` handles restored pipelines. A queued, never-started 
pipeline therefore holds the executor object but no threads.
   - The accurate bound today is "up to two worker threads per pipeline whose 
scheduler has actually been used", not "two per submitted job".
   
   ## Motivation
   
   An operator consolidating many small streaming jobs onto one Zeta member 
should not pay two sleeping OS threads per checkpoint-enabled pipeline. The 
cost is not CPU; it is thread-count pressure, stack memory, and the scheduler 
noise of a large pool of idle timers, all of which grow with a number the 
operator controls only indirectly.
   
   After this change, the steady-state timer cost of a member becomes a small 
constant independent of pipeline count, while per-pipeline cancellation 
semantics and the single-pending-checkpoint invariant stay exactly as they are 
today.
   
   This is a scalability change to thread ownership. It is not a regression 
report and not a change to checkpoint semantics.
   
   ## Goals
   
   - Replace the per-coordinator scheduled pool with a member-wide timer whose 
thread count does not grow with the number of active pipelines.
   - Keep cancellation scoped per pipeline, keyed by `(jobId, pipelineId)`, so 
one pipeline's cleanup cannot touch another's outstanding timers.
   - Isolate blocking checkpoint or RPC work so that a slow pipeline delays 
neither another pipeline's timers nor its own next tick.
   - Preserve the single-pending-checkpoint guard and the existing locking 
discipline in the trigger path.
   - Keep the coordinator usable across a master-failover reset without 
rebuilding any thread pool.
   - Introduce no new `Option`, no default change, and no checkpoint, 
savepoint, or protocol format change.
   
   ## Non-Goals
   
   - Changing checkpoint triggering semantics, barrier alignment, checkpoint 
storage, or savepoint behavior.
   - Bounding or reshaping the coordinator executor. That is separate work, 
tracked in #12117 and #12121, and this proposal deliberately stays independent 
of it.
   - Making the checkpoint interval, timeout, or any scheduling parameter 
user-tunable beyond what exists today.
   - Flink or Spark translation layers. This is a Zeta engine change only.
   - Cross-member or cluster-wide scheduling. The unit of sharing is one member.
   
   ## Design
   
   ### 1. Ownership
   
   One `SharedCheckpointScheduler` per member, owned by 
`SeaTunnelEngineContext` and closed with it. It is deliberately not a static 
singleton, because several members run inside one JVM during unit and E2E 
tests, and a static instance would let one member's shutdown tear down 
another's timers.
   
   Each `CheckpointCoordinator` obtains a `PipelineCheckpointScheduler` lease 
instead of constructing a pool. The lease is identified by `(jobId, 
pipelineId)` and holds only its own outstanding tasks. Nothing is registered 
back on the shared object, so dropping a lease when a job ends leaks nothing.
   
   `cleanPendingCheckpoint` cancels this pipeline's outstanding tasks rather 
than shutting down a pool. The lease stays usable afterwards, so a coordinator 
restored after a master-failover reset can schedule again immediately.
   
   ```mermaid
   graph TD
       subgraph Today["Today: pool per pipeline"]
           C1[Coordinator p1] --> P1["pool: 2 threads"]
           C2[Coordinator p2] --> P2["pool: 2 threads"]
           CN[Coordinator pN] --> PN["pool: 2 threads"]
       end
   
       subgraph Proposed["Proposed: leases on member threads"]
           L1[Coordinator p1] --> S["SharedCheckpointScheduler<br/>owned by 
SeaTunnelEngineContext"]
           L2[Coordinator p2] --> S
           LN[Coordinator pN] --> S
       end
   ```
   
   ### 2. Split timing from execution
   
   Timing and execution are handled by two different pools, and this is the 
part that most deserves review.
   
   A fixed, small timer pool does nothing but hand a due task to a dispatch 
pool. The checkpoint body itself never runs on a timer thread. The hand-off is 
a constant-time, non-blocking `execute` onto a cached pool, which reuses an 
idle thread or creates one rather than queueing.
   
   The consequence is that a checkpoint that blocks on IO or a slow RPC 
occupies a dispatch thread, not a timer thread. Every other pipeline's next 
trigger still fires on schedule, and so does the slow pipeline's own next tick.
   
   This is stricter than today's behavior, not merely equivalent. Today a 
coordinator's two threads serve both timing and bodies, so a blocked body 
consumes one of that pipeline's own timer threads, and a blocked trigger 
together with a blocked watchdog can occupy both.
   
   ```mermaid
   sequenceDiagram
       autonumber
       participant Lease as PipelineCheckpointScheduler
       participant Timer as Timer pool (fixed, small)
       participant Disp as Dispatch pool (elastic)
       participant Coord as CheckpointCoordinator body
   
       Lease->>Timer: schedule(task, delay)
       Note over Timer: sleeps until due
       Timer->>Disp: execute(task)
       Note over Timer: returns immediately,<br/>free for the next pipeline
       Disp->>Coord: run trigger or timeout watchdog
       Note over Coord: may block on RPC,<br/>holds a dispatch thread only
   ```
   
   ### 3. Cancellation across two pools
   
   Because a task crosses two pools, cancellation has to work both before the 
timer fires and after the body has been handed to the dispatch pool. The task 
therefore re-checks its own cancellation flag at run time rather than relying 
on the timer future, which is already complete once the body is running. The 
timer uses `setRemoveOnCancelPolicy(true)`, since coordinators cancel the 
timeout watchdog on every acknowledged checkpoint and cancelled entries would 
otherwise sit in the shared queue until their delay elapsed.
   
   ### 4. Invariants
   
   The single-pending-checkpoint guard is untouched. The trigger path stays 
inside `synchronized (lock)` and the pending-counter guard still short-circuits 
a redundant trigger. Coordinators already had two threads able to enter that 
path concurrently, so sharing the timer introduces no new race: additional 
dispatch threads contend on the same lock and bail through the same guard.
   
   ### 5. Open question: bounding the dispatch pool
   
   Issue #12122 asks for a "shared, bounded scheduler". The timer pool is 
bounded. The dispatch pool, as currently prototyped, is not: it is elastic and 
reaps idle threads after the cached pool's keep-alive.
   
   The tradeoff, stated plainly:
   
   - Unbounded dispatch means that if N pipelines block simultaneously, up to N 
dispatch threads exist for as long as they block. That worst case is no worse 
than today's fixed two-per-pipeline, and the steady state is far better, but it 
is not a hard bound.
   - A bounded pool over a `SynchronousQueue` rejects when saturated. A 
rejected task here is a silently dropped checkpoint trigger or a silently 
dropped timeout watchdog, which is a correctness problem rather than a resource 
problem.
   
   Making the dispatch pool bounded therefore requires choosing a rejection 
policy, and the plausible options each have costs: caller-runs would push 
checkpoint work back onto a timer thread and reintroduce exactly the 
head-of-line blocking this design removes; an unbounded queue trades thread 
growth for latency and memory growth; a bounded queue with an explicit failure 
path needs agreement on what a member should do when it cannot trigger a 
checkpoint on time.
   
   This is the main question I would like the list to settle, because it 
decides whether the implementation ships as-is or grows a rejection contract.
   
   A secondary question: the timer pool size is currently a constant rather 
than an `Option`. Adding one would create a user-facing config contract for 
something that arguably should not need tuning, and would require matching 
`docs/en` and `docs/zh` updates. I have no strong view and will follow the list.
   
   ## Compatibility and Migration
   
   No incompatible change is proposed.
   
   - No new or renamed `Option`, and no default value change.
   - No checkpoint, savepoint, or wire format change; no SPI or public API 
change.
   - No entry in `docs/en/introduction/concepts/incompatible-changes.md` is 
required.
   - The only externally observable difference is thread naming. 
Per-coordinator threads named `checkpoint-coordinator-<pipelineId>/<jobId>` are 
replaced by member-wide pools. Anything scraping thread names, including the 
E2E thread-leak check, needs updating; the candidate PR updates that check.
   
   ## Test Plan
   
   - Unit tests asserting that the member timer thread count stays flat as 
pipeline count grows, that cancellation is scoped to one lease and leaves other 
leases untouched, that a lease remains usable after `cancelAll`, and that 
cancellation works both before the timer fires and after the body is dispatched.
   - Existing `CheckpointCoordinatorTest` must pass unchanged in behavior, 
covering restore, savepoint, and cleanup paths.
   - E2E: the Zeta engine suites, including checkpoint restore and savepoint 
scenarios, plus the container thread-leak check with the new pool names 
accounted for.
   - Scalability check per #12122: benchmark active checkpoint-enabled 
pipelines separately from initialized or queued jobs, and confirm 500 active 
pipelines show a flat timer-thread count.
   - Not yet covered and worth adding if the list wants it: a test that 
deliberately blocks one pipeline's checkpoint body and asserts another 
pipeline's trigger still fires within its interval.
   
   ## Alternatives Considered
   
   1. **Dispatch checkpoint work onto the existing coordinator executor.** 
Rejected. That executor is a bounded `ThreadPoolExecutor` over a 
`SynchronousQueue` with a rejection handler, so a saturated coordinator pool 
could silently drop a checkpoint trigger or timeout watchdog. It would also 
couple this change to the executor-bounding work in #12117 and #12121.
   
   2. **One shared pool for both timing and bodies.** Rejected. A blocking body 
would consume a shared timer thread and delay every other pipeline's tick, 
which is the exact failure mode #12122 asks to avoid.
   
   3. **Keep per-coordinator pools but shrink them to one thread.** Halves the 
cost but does not change the linear growth, and makes self-interference worse, 
since a single thread serves both the trigger and the watchdog.
   
   4. **A static singleton scheduler.** Rejected. Several members run in one 
JVM in tests and E2E, so one member's shutdown would tear down another's timers.
   
   5. **Leave as-is.** Acceptable for small deployments and is the status quo; 
it simply does not address the consolidation case in #12122.
   
   ## References
   
   - Issue [#12122](https://github.com/apache/seatunnel/issues/12122) - Share 
checkpoint trigger scheduling across active pipelines
   - PR [#12165](https://github.com/apache/seatunnel/pull/12165) - candidate 
implementation and review thread
   - Related executor-bounding work: 
[#12117](https://github.com/apache/seatunnel/issues/12117), 
[#12121](https://github.com/apache/seatunnel/issues/12121)
   


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