Vivek1106-04 opened a new pull request, #12277:
URL: https://github.com/apache/seatunnel/pull/12277

   ### Purpose of this pull request
   
   Closes #12121.
   
   When thread sharing is enabled (`task_execution_thread_share_mode` is `ALL` 
or `PART`), many tasks share one cooperative worker thread. 
`TaskCallTimer.timeoutAct` promotes a worker whenever a task call outlives the 
50 ms call timer: the worker becomes exclusive to that slow task 
(`CooperativeTaskWorker.exclusiveTaskTracker`) and 
`RunBusWorkSupplier.runNewBusWork(false)` starts a replacement for the shared 
queue.
   
   Every promotion therefore costs one thread, and nothing bounded it: 
`executorService` is a cached pool and `threadShareTaskQueue` is unbounded, so 
the worker thread count grew with the number of slow cooperative calls rather 
than with the number of slots.
   
   This PR makes that growth an explicit, observable admission decision.
   
   **What it adds**
   
   * `CooperativeWorkerBudget` - a promotion budget with a global limit and a 
per job limit, plus counters for promotions and denials. Budget is acquired 
before a promotion and released when the promoted worker is done with its task, 
not when a timer fires.
   * `TaskCallTimer` no longer promotes unconditionally. When the budget denies 
a promotion the timeout is not dropped: the slow task stays where it is and the 
promotion is retried with a bounded exponential backoff (up to 1 s), so a 
released budget is picked up by the tasks that are waiting for it.
   * Reserved readiness capacity: when the denied worker is the last one 
serving the shared queue, a replacement worker is still started. An exhausted 
promotion budget can therefore never prevent source, sink, or coordinator tasks 
from starting.
   * `RunBusWorkSupplier.tryPromoteCooperativeWorker` is the single place that 
makes the decision, so the job id, the budget, and the replacement worker stay 
together.
   * The worker thread pool status log now reports `sharedCooperativeWorkers`, 
`promotedCooperativeWorkers`, `totalCooperativePromotions`, and 
`deniedCooperativePromotions`; denied promotions are logged with an explicit 
`BUDGET_EXHAUSTED` reason.
   
   **Configuration**
   
   | option | default | meaning |
   |---|---|---|
   | `max-promoted-cooperative-workers` | `0` | global limit of promoted 
cooperative workers on a worker node; `0` is unlimited |
   | `max-promoted-cooperative-workers-per-job` | `0` | limit per job, so one 
job cannot consume the whole budget; `0` is unlimited |
   
   Both default to `0`, so an upgraded cluster behaves exactly as before until 
an operator opts into a bound. Documented in the English and Chinese hybrid and 
separated deployment guides.
   
   ### Does this PR introduce _any_ user-facing change?
   
   Yes, additive only: two new optional worker options, documented in `docs/en` 
and `docs/zh`. Defaults keep the existing behavior (unlimited promotions), and 
no existing option or default is changed.
   
   ### How was this patch tested?
   
   * `CooperativeWorkerBudgetTest` (new) - unlimited and negative limits, 
global limit, per job limit including the rollback of the global reservation 
when the per job limit denies, release and reuse, release of an unknown job, 
and a 32 thread concurrent acquire that must admit exactly the limit.
   * `TaskExecutionServiceCooperativeBudgetTest` (new) - deploys 8 cooperative 
tasks whose every call takes 300 ms against a node configured with a global 
limit of 2 and a per job limit of 1 
(`seatunnel_cooperative_worker_budget.yaml`). It asserts that promotions are 
denied once the job holds its single promoted worker, that promoted workers 
never exceed the budget, that the workers serving the shared queue stay 
bounded, that every task keeps being called while the budget is exhausted (no 
readiness deadlock), that the task group still finishes, and that the budget is 
returned afterwards.
   * `YamlSeaTunnelConfigParserTest` - parses both new options and covers the 
defaults and the rejection of negative values.
   * Regression: `TaskExecutionServiceTest` (15 tests), 
`seatunnel-engine-common` tests, and `ImportClassCheckTest` all pass; `./mvnw 
-pl 
seatunnel-engine/seatunnel-engine-common,seatunnel-engine/seatunnel-engine-server
 -DskipTests verify` passes with spotless applied.
   
   ### Check list
   
   * [x] If any new Jar binary package adding in your PR, please add License 
Notice according [New License 
Guide](https://github.com/apache/seatunnel/blob/dev/docs/en/developer/new-license.md)
 - no new dependency
   * [x] If necessary, please update the documentation to describe the new 
feature. https://github.com/apache/seatunnel/tree/dev/docs
   * [x] If necessary, please update `incompatible-changes.md` to describe the 
incompatibility caused by this PR - not needed, defaults keep the previous 
behavior
   * [x] If you are contributing the connector code, please check that the 
following files are updated - not a connector change
   
   🤖 Generated with [Claude Code](https://claude.com/claude-code)
   
   https://claude.ai/code/session_01HoCQS4kt6SNUeo6jMXK8d7
   


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