Skip to content

Worker Pool 与 Scheduler

Mini Ray 用 ThreadPoolExecutor 模拟 Worker Pool,用 ResourceScheduler 模拟 Ray 的资源调度。

任务提交

RemoteFunction.remote() 最终会调用:

python
Runtime.submit_task(fn, args, kwargs, options)

执行过程:

注意:ObjectRef 在任务真正执行前就返回给调用方。

资源调度器

Mini Ray 的资源调度器维护两个字典:

python
total = {"CPU": 4, "GPU": 1}
available = {"CPU": 4, "GPU": 1}

执行前尝试获取资源:

python
def acquire(self, request):
    while not self._can_fit(request):
        self._cond.wait()
    for name, amount in request.items():
        self._available[name] -= amount

执行后释放:

python
def release(self, request):
    for name, amount in request.items():
        self._available[name] += amount
    self._cond.notify_all()

这就是最小可用的资源调度模型。

为什么线程数和 CPU 资源不是一回事

Mini Ray 的 executor 线程数设置为 num_workers * 4,但 CPU 资源由 ResourceScheduler 控制。

原因是有些线程可能在等待依赖对象;如果线程数等于 CPU 资源数,很容易出现“线程都在等对象,没有线程去生产对象”的教学版死锁。

真实 Ray 中,这部分由 Core Worker、Raylet、依赖解析和 Worker lease 共同处理。

运行资源 Demo

bash
uv run python examples/mini_ray_runtime/demos/05_scheduler_resources.py

其中:

python
ray.init(num_workers=4, resources={"CPU": 4, "GPU": 1})

@ray.remote(num_cpus=1, num_gpus=1)
def gpu_task(i):
    time.sleep(0.2)
    return f"gpu-{i}"

即使 CPU 足够,GPU 只有 1 个,GPU 任务也会串行通过调度器。

调度器还能怎么扩展

你可以继续实现:

  • 优先级队列。
  • FIFO 任务队列。
  • 按对象位置做本地性调度。
  • Placement Group。
  • 任务取消。
  • 超时和资源抢占。

但不要一开始就做这些。先保证 acquire -> run -> release 这条主线可靠。

小练习

metrics() 中增加 queued_tasks。提示:当前实现没有显式任务队列,你需要在 submit_task 时计数,在获取资源后减少。

面向学习目的的 Ray Core 中文导读与 Mini Ray 机制预览。