This repository has been archived by the owner on Feb 22, 2022. It is now read-only.
[stable/dask] ability to set worker --nthreads and --nprocs independently of cpu resource limit #18708
Labels
lifecycle/stale
Denotes an issue or PR has remained open with no activity and has become stale.
A dask worker node should be able to have --nthreads and --nprocs set via a user-provided helm chart, and have those values be independent of worker cpu resources. One use case is high I/O workloads that benefit from more threads than cores.
Currently --nthreads is set to the worker cpu resource limit.
Aside from limiting dask helm deployment options, this limitation is problematic because it leads users to believe that there needs to be a 1-1 relationship between cpus and threads, which is not true.
A solution is to have a new option in the helm chart:
Other worker command line options could be added in this way.
The text was updated successfully, but these errors were encountered: