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]
