Docker部署Dask后,如何限制map函数的并行Worker数量?
解决方案
你可以通过以下几种方式限制该map任务的并行执行数量,同时保留20个Worker容器供其他任务使用:
方法1:使用client.map的concurrency_limit参数
Dask的Client.map方法支持concurrency_limit参数,可直接指定该任务允许同时运行的最大并发数。这种方式既能限制当前map任务的并行度,又不影响其他任务占用剩余的Worker资源。
修改代码如下:
worker_args = .... # 包含100个Worker参数的数组 # 限制该map任务最多同时运行5个Worker futures = client.map(function_in_worker, worker_args, concurrency_limit=5) worker_responses = client.gather(futures)
方法2:分批提交任务并控制并发
如果需要更精细的执行顺序控制,可以手动将任务分批次提交,始终保持同时运行的任务数不超过5个:
from dask.distributed import as_completed worker_args = .... # 包含100个Worker参数的数组 batch_size = 5 worker_responses = [] # 提交第一批任务 futures = [client.submit(function_in_worker, arg) for arg in worker_args[:batch_size]] # 循环处理完成的任务并补充新任务 for future in as_completed(futures): worker_responses.append(future.result()) # 还有未提交的任务时,提交下一个 if worker_args: next_arg = worker_args.pop(0) new_future = client.submit(function_in_worker, next_arg) futures.append(new_future)
方法3:用delayed构建任务图并指定计算并发数
若愿意调整代码结构,可通过delayed装饰函数构建任务图,再通过compute方法指定num_workers参数限制并发:
from dask import delayed worker_args = .... # 包含100个Worker参数的数组 # 构建延迟任务列表 delayed_tasks = [delayed(function_in_worker)(arg) for arg in worker_args] # 限制计算时最多使用5个Worker worker_responses = client.compute(delayed_tasks, num_workers=5)
内容的提问来源于stack exchange,提问作者ps0604
相关产品推荐
相关产品推荐

