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]
