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 时计数,在获取资源后减少。