调度与资源模型
Ray 调度器要回答的核心问题是:这个 Task 或 Actor 应该在哪个节点、哪个 Worker 上运行?
只找一个空闲 CPU 不够。调度器还要考虑资源声明、对象位置、节点健康状态、Placement Group 和队列压力。
逻辑资源
Ray 的资源是逻辑资源,不是操作系统级别的硬限制。
@ray.remote(num_cpus=2, num_gpus=1)
def train(config):
...含义是:调度这个任务时需要一个拥有至少 2 个 CPU 资源和 1 个 GPU 资源的节点。
注意两点:
num_cpus不会强制限制你的函数只能开两个线程。- 如果底层库会自动使用很多线程,你需要自己配置,例如
OMP_NUM_THREADS。
自定义资源
自定义资源适合表达特殊约束:
ray start --resources='{"ssd": 2, "license": 1}'@ray.remote(resources={"ssd": 1})
def load_from_fast_disk(path):
...这不是物理隔离,而是调度标签和配额。
调度决策
一次任务调度通常会考虑:
| 因素 | 为什么重要 |
|---|---|
| CPU / GPU / 自定义资源 | 节点是否有足够资源 |
| 对象位置 | 避免把大对象拉到远端 |
| Worker 可用性 | 是否需要启动新 Worker |
| Placement Group | 一组任务是否需要共同布局 |
| 节点健康状态 | 避免调度到不可用节点 |
| 队列压力 | 避免局部拥塞 |
数据局部性
假设任务需要读取一个 10GB 对象。如果对象已经在 Node A,本地运行通常比调度到 Node B 再传输更好。
Ray 不会只看 CPU 空闲,还会考虑对象位置和传输成本。
Placement Group
有些任务必须成组调度。比如分布式训练需要 4 个 Worker 同时占用 GPU,如果只启动 2 个,整个训练也跑不起来。
Placement Group 让你把资源请求打包:
from ray.util.placement_group import placement_group
pg = placement_group([
{"GPU": 1},
{"GPU": 1},
{"GPU": 1},
{"GPU": 1},
], strategy="PACK")
ray.get(pg.ready())常见策略:
| 策略 | 含义 |
|---|---|
PACK | 尽量放在少数节点,减少通信 |
SPREAD | 尽量分散,提高容错或隔离 |
STRICT_PACK | 必须放在一个节点 |
STRICT_SPREAD | 必须分散到不同节点 |
反压与资源饥饿
如果 Driver 一次提交几万个任务,任务队列会很长,占用大量内存。Ray 能排队,但业务层也需要控制并发量,避免“提交太多、内存爆掉”。
常见做法:
- 用
ray.wait控制同时在途的任务数量。 - 把小任务合成 batch。
- 用 Actor 做生产者/消费者模式。
- 对 GPU 等稀缺资源显式声明,让调度器自动串行化。
示例:限制最多 32 个 pending 任务。
pending = []
for item in items:
pending.append(work.remote(item))
if len(pending) >= 32:
ready, pending = ray.wait(pending, num_returns=1)
consume(ray.get(ready[0]))
for ref in pending:
consume(ray.get(ref))Mini Ray 的调度器
Mini Ray 的 ResourceScheduler 是一个最小版本:
- 维护
total和available。 - 任务提交后先解析依赖。
- 执行前
acquire(resources)。 - 执行结束后
release(resources)。
@ray.remote(num_cpus=1, num_gpus=1)
def gpu_task(i):
return i
ray.init(num_workers=4, resources={"CPU": 4, "GPU": 1})
refs = [gpu_task.remote(i) for i in range(3)]虽然有 4 个 CPU,但只有 1 个 GPU,所以这些任务会串行占用 GPU 资源。
为什么 Mini Ray 不做全局调度
真实 Ray 面对多节点,需要:
- 节点间资源视图。
- 任务转发和 lease。
- 对象位置索引。
- Worker 生命周期管理。
- 节点故障处理。
Mini Ray 只有一个进程,所以它的调度器只解决“本地资源并发控制”。这正好适合作为第一版:你先理解资源计数和队列,再去看 Raylet 的多节点逻辑。
小练习
运行:
uv run python examples/mini_ray_runtime/demos/05_scheduler_resources.py把 resources={"CPU": 4, "GPU": 1} 改成 {"CPU": 4, "GPU": 2},观察总耗时变化。再把 @ray.remote(num_cpus=1, num_gpus=1) 改成 num_gpus=0,思考为什么调度行为会改变。