如何在Dask中当任务批次优先级相同时随机选择批次执行下一个任务?
如何在Dask中当任务批次优先级相同时随机选择批次执行下一个任务?
首先得说,你找对方向了——调度器插件就是实现这种自定义任务调度逻辑的最佳途径,不用去动Dask的核心源码,灵活又好维护。默认的FIFO/LIFO逻辑确实不会考虑这种“同优先级下随机选批次”的需求,咱们直接用插件来改写任务的挑选逻辑就行。
我给你写一个具体的插件实现,你可以直接拿去用,核心就是重写调度器里挑选下一个要执行的任务的逻辑:
from dask.distributed import SchedulerPlugin import random class RandomBatchScheduler(SchedulerPlugin): def __init__(self): super().__init__() def choose_task(self, scheduler, worker): # 先收集所有处于可运行状态的任务批次(也就是有等待执行任务的批次) ready_batches = [] for client_key, client_state in scheduler.clients.items(): # 检查这个客户端的任务队列里有没有可运行的任务 if client_state.ready: ready_batches.append(client_state) if not ready_batches: # 没有可运行任务,返回None让调度器按默认逻辑走 return None # 随机选一个批次 selected_batch = random.choice(ready_batches) # 从这个批次里拿出下一个任务(这里保持批次内的FIFO逻辑,你也可以改成随机,看你需求) return selected_batch.ready.popleft()
然后你需要在启动调度器的时候注册这个插件,或者在客户端连接后注册(如果是用现有调度器的话):
如果是自己启动调度器:
from dask.distributed import Scheduler, Client # 启动调度器并注册插件 scheduler = Scheduler() scheduler.add_plugin(RandomBatchScheduler()) scheduler.start("tcp://0.0.0.0:8786") # 客户端连接 client = Client("tcp://0.0.0.0:8786")
如果是连接已经运行的调度器:
from dask.distributed import Client client = Client("tcp://your-scheduler-ip:8786") client.register_plugin(RandomBatchScheduler(), scheduler=True)
逻辑解释:
choose_task方法是调度器插件的核心钩子,每次worker有空位时,调度器都会调用这个方法来选下一个任务- 我们先遍历所有客户端的任务状态,把有可运行任务的批次(也就是你说的两个不同用户的任务批次)收集起来
- 用
random.choice随机选一个批次,然后从这个批次的就绪队列里拿出第一个任务(这里保持批次内的FIFO逻辑,如果你需要批次内也随机,把popleft()改成pop(random.randint(0, len(selected_batch.ready)-1))就行) - 如果没有就绪任务,就返回None,让调度器按默认逻辑处理
为什么不用改源码?
直接改Dask源码的话,你得去修改scheduler.py里的choose_task或者相关的任务队列管理逻辑,但这样做的问题是,每次Dask版本更新你都得重新改,维护成本太高。用插件的话,完全是增量式的自定义,不影响核心逻辑,也方便升级。
另外你提到的fifo_timeout确实只是控制FIFO和LIFO的切换,和随机选批次完全没关系,所以这个参数解决不了你的问题。
最后提醒一下,这个插件是按客户端来区分批次的,如果你说的“批次”是同一个客户端下的不同任务组,那你需要稍微调整一下逻辑——比如给每个任务加一个自定义的批次标签,然后在choose_task里按标签分组,再随机选组。不过看你的描述,两个批次是不同用户不同客户端提交的,所以上面的代码刚好适配你的场景。
备注:内容来源于stack exchange,提问作者Z4NG
相关产品推荐
相关产品推荐

