询问dask.distributed中负责向Worker发送任务的代码位置
Dask Distributed调度器向Worker发送任务的代码定位指引
- 任务分配的触发时机:当Worker上报空闲(
worker_ready事件)、新任务加入等待队列,或是Worker完成任务(task_finished事件处理后),调度器会启动任务分配流程。 - 核心分配入口:
Scheduler.assign方法(位于distributed/scheduler.py)是调度器分配任务的核心逻辑入口,它会依据Worker负载、数据本地性等规则筛选合适的目标Worker。 - 实际发送任务的方法:
Scheduler.send_to_worker负责将任务指令(如compute类型的消息)发送给指定Worker。在assign方法确定目标Worker后,会调用这个方法完成最终的任务发送动作。 - 代码细节参考:
- 在
assign方法内,会调用choose_worker选择最优Worker,生成待发送的任务列表; - 随后通过
send_to_worker把任务打包成消息,通过底层通信通道(如Tornado的IOLoop)发送给Worker; - 搜索
send_to_worker的定义,能看到它如何构建任务消息并调用worker.send()完成实际发送。
- 在
内容的提问来源于stack exchange,提问作者Nipun Talukdar
相关产品推荐
相关产品推荐

