You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Ubuntu集群中Dask任务无法在Worker间并行运行问题咨询

解决Dask任务未调度到空闲Worker的问题

我来帮你揪出这个Dask任务调度的问题!你遇到的情况——第一个任务占着Worker,第二个任务却没去空闲的第二个Worker——大概率是这几个原因导致的,咱们一步步排查:

1. 任务粒度太小,Dask默认打包调度

Dask为了减少调度开销,会尽量把小任务打包到同一个Worker上执行。如果你的my_task执行时间很短(比如几毫秒就跑完),调度器会觉得没必要折腾到另一个Worker,直接塞给第一个Worker处理。

解决办法:

  • 给my_task加个延迟,比如time.sleep(10),让任务执行时间足够长,这样第一个Worker被占满后,调度器就会把第二个任务分配到空闲的Worker。
  • 或者调整客户端的任务打包参数,比如client.submit(my_task, arg, batch_size=1),强制每个任务单独调度。

2. Worker资源配置没限制,第一个Worker还有空闲线程

默认情况下,Dask Worker会使用机器上的所有CPU核作为线程数。如果你的Worker机器有多核,第一个Worker可能还有空闲线程,调度器自然会把第二个任务发过去,而不是用第二个Worker。

解决办法:
启动Worker的时候,明确指定每个Worker只使用1个线程,这样第一个任务会占满整个Worker:

# 在两台Worker机器上分别执行
dask-worker tcp://你的调度器IP:8786 --nthreads 1 --nworkers 1

3. 先确认Worker是否真的正常连接

有时候可能Worker没连上调度器,你以为有两个Worker,实际只有一个在运行。在客户端代码里加一行检查:

print("当前在线Worker:", client.ncores())

如果输出是类似{'tcp://worker1:xxxx': 1, 'tcp://worker2:yyyy': 1},说明两个Worker都在线;如果只有一个,那得先排查Worker的连接问题。

调试用的示例代码

我给你改了下客户端代码,加了机器名打印和延迟,方便你验证:

from dask.distributed import Client
import time
import os
import random

def my_task(arg):
    # 打印任务运行的机器名,方便确认是哪个Worker在处理
    print(f"任务 {arg} 正在 {os.uname()[1]} 上运行")
    time.sleep(10)  # 延长任务执行时间,确保Worker被占用
    return random.randint(0, 100)

if __name__ == "__main__":
    # 替换成你的调度器IP和端口
    client = Client("tcp://192.168.x.x:8786")
    print("已连接调度器,在线Worker资源:", client.ncores())
    
    # 提交第一个任务
    future1 = client.submit(my_task, 1)
    print("已提交任务1")
    
    # 等2秒,确保第一个任务开始执行
    time.sleep(2)
    
    # 提交第二个任务
    future2 = client.submit(my_task, 2)
    print("已提交任务2")
    
    # 获取结果
    res1 = future1.result()
    res2 = future2.result()
    print(f"任务1结果: {res1}, 任务2结果: {res2}")
    
    client.close()

先按上面的步骤排查,大概率能解决问题~

内容的提问来源于stack exchange,提问作者TheCodeCache

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.25 07:17:07