前置补充:Raylet 和 GCS
Raylet(运行于集群的每个节点上)
- 主要功能包括:
- 本地资源管理 :监控并分配所在节点的 CPU、GPU 及内存资源
- 工作进程管理 :维护一个工作进程池(Worker Pool),负责任务执行进程的分配、回收与异常处理
- 本地对象存储 :管理节点上的共享内存存储(Object Store),处理本地数据的读写与内存回收
- 信息上报 :定期向 GCS 发送心跳,上报本节点的资源状态、任务执行情况及对象存储位置
- Raylet 是无状态的
- 自身状态存储在 GCS 中,故障后可从 GCS 恢复
GCS(运行于集群的头节点)
- GCS(Global Control Store)是 Ray 的集中式元数据管理组件
- 主要功能包括:
- 元数据存储 :作为中心化键值存储(后端常为 Redis 或 RocksDB),保存全局状态
- 包括节点拓扑、Actor 位置、分布式对象的位置映射、Job 及 Task 的执行记录
- 集群协调 :处理节点注册,通过心跳检测节点存活状态,触发故障节点的标记与恢复流程
- 全局服务 :提供 Actor 全局命名与定位、Placement Group 管理、分布式锁等原子化服务
- 元数据存储 :作为中心化键值存储(后端常为 Redis 或 RocksDB),保存全局状态
- GCS 是有状态的
- 需配置高可用后端以支持容错
Raylet 与 GCS 的协作
- 状态汇总与调度 :各 Raylet 上报状态至 GCS,GCS 构建全局资源视图
- 调度器查询该视图选定目标节点,并向该节点的 Raylet 下发任务执行指令
- 数据定位与传输 :工作进程读取远程对象时,本地 Raylet 向 GCS 查询该对象的存储节点列表,随后通过点对点网络直接从目标节点的 Raylet 关联存储中拉取数据
- 故障处理 :GCS 检测到某 Raylet 失联后,停止向该节点分配任务,并触发其上任务和 Actor 的重调度
- 若 GCS 自身故障,各 Raylet 依据高可用配置尝试重连或选举新实例
资源定义(节点侧容量声明)
- Ray 的资源模型基于键值对(Key-Value)
- Key 为资源名称(字符串)
- Value 为 容量(浮点数)
Ray 资源分类
- 内置资源(Built-in) :
CPU、GPU、memory、object_store_memory- 这四个名称被 Ray 调度器(GCS,Global Control Store)硬编码识别
- 自定义资源(Custom) :用户任意命名的资源,例如
inference_license、cache_slots等自定义名称- 注:自定义资源名称不能与内置资源(CPU/GPU/memory/object_store_memory)重名,否则会被当作内置资源处理或报错
节点启动时的容量声明
在启动 Ray 节点进程时,通过命令行参数固定该节点的资源总容量,这些数值注册到全局状态存储(即存储到 GCS 中)后不可动态修改
1
2# 声明该节点拥有 8 个 CPU 单位、4 个 GPU 单位、1 个自定义资源单位
ray start --num-cpus=8 --num-gpus=4 --resources='{"inference_license": 1}'- 这里的声明是逻辑资源,所以
--num-cpus=1000也可以,完全可以超过真实 cpu 核数- 但这个节点总资源是用于限制 Ray 提交太多任务到同一个节点的,声明过多逻辑资源容易造成资源抢占甚至崩溃
--num-cpus默认值物理 CPU 核心数(占用全部物理 CPU 核心),对于头节点(包含--head参数)的节点,建议设置为 0
- 这里的声明是逻辑资源,所以
资源请求与限制(Task/Actor 侧语法)
- 用户在提交计算任务时,通过
ray.remote装饰器或.options()方法声明资源需求量- 提交任务时的需求量是硬性约束 ,调度器仅在有节点剩余容量满足该数值时才会进行调度(对应前面节点的逻辑资源申明)
- 调用时的任务的
.options()参数可以覆盖定义时的ray.remote参数
Ray 中的任务资源申请
Ray 资源是逻辑资源,非物理隔离,更多是用于调度时的准入控制
- Ray 不会将一个物理 CPU 核心绑定并独占给这个任务
- Ray 甚至不会阻止一个
num_cpus=1的任务在内部启动多个线程num_cpus=1下内部也可以使用多个线程使用多个物理 CPU 核心- 实际的 CPU 使用由操作系统调度
Task :默认需求为
num_cpus=1- 若不显式设置,每个提交的任务都会消耗 1 单位 CPU 逻辑配额
Actor :默认需求为
num_cpus=0Task(无状态函数)资源示例:
1
2
3
4
5
6
7
8
9
10
11import ray
ray.init()
# 定义阶段:装饰器作用于函数
def my_function(x):
return x + 1
# 提交阶段:使用 .options() 动态覆盖资源需求
# .remote() 执行后返回 ObjectRef(对象引用)
task_ref = my_function.options(num_cpus=0.5, resources={"inference_license": 1}).remote(10)my_function是 Ray 封装的远程可调用对象my_function.options(...)生成一个携带新资源配置的副本.remote(10)将任务下发到集群,返回占位符ObjectRef
Actor(有状态类)资源申请示例
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17import ray
ray.init()
# 定义阶段:装饰器作用于类
class MyActor:
def __init__(self, val):
self.val = val
def get(self):
return self.val
# 创建阶段:使用 .options() 动态覆盖资源需求
# .remote(init_args) 在远端创建实例,返回 ActorHandle(句柄)
actor_handle = MyActor.options(num_cpus=0.5, resources={"inference_license": 1}).remote(100)
# 通过句柄调用方法,返回 ObjectRef
result_ref = actor_handle.get.remote()MyActor是Ray封装的远程Actor构造器MyActor.options(...)生成一个携带新资源配置的构造器副本.remote(100)在远端节点实例化对象,返回ActorHandle(非本地类实例)- 设置
num_cpus=0.5表明该任务声明需要 0.5 个单位的逻辑 CPU 配额- 调度准入时仅占用 0.5 cpu
- 物理执行时理论上可以多用
- 环境变量设置(软性建议)
- Ray 为执行该任务的工作进程设置环境变量
OMP_NUM_THREADS,规则为将num_cpus向上取整 :0.5向上取整为1,故OMP_NUM_THREADS=1- Ray 确实会设置
OMP_NUM_THREADS,但并非所有库都遵守(如NumPy默认不限制)
- 科学计算库(OpenBLAS、MKL、PyTorch)在初始化时读取此变量作为线程池上限,从而软性建议该任务减少内部并行度,避免挤占其他任务的物理 CPU 时间
- Ray 为执行该任务的工作进程设置环境变量
任务提交后的调度器的资源分配流程
- 当任务提交后,Ray 的 GCS 调度器执行以下确定性流程:
- 步骤一:可行性检查(Feasibility)
- 遍历所有存活节点,剔除两类节点:
- 1)资源类型缺失 :例如在无 GPU 节点上请求
num_gpus=1(注意:这里是提交任务时,不是初始化节点进程时) - 2)剩余容量不足 :节点的某种资源剩余量
<任务请求量
- 1)资源类型缺失 :例如在无 GPU 节点上请求
- 遍历所有存活节点,剔除两类节点:
- 步骤二:节点排序(Scoring)
- 默认调度策略(
DEFAULT)对候选节点进行综合打分:- 数据本地性优先 :若任务的输入数据(通过
ray.put存储)已存在于某节点的对象存储中,该节点得分大幅提升 - 负载均衡优先 :计算每个节点资源利用率(已占用/总容量)
- 节点剩余可用资源越多,得分越高
- 数据本地性优先 :若任务的输入数据(通过
- 默认调度策略(
- 步骤三:Top-K 随机选取
- 为避免所有任务涌入得分最高的单一节点,调度器从 排序前 K 个 节点中随机抽取一个
- K 值默认为集群总节点数的 20%(由环境变量
RAY_SCHEDULER_TOP_K_AWARE=1启用,RAY_SCHEDULER_TOP_K_AWARE_FRACTION控制)
高级资源分配机制
自定义资源(数值型准入)
- 自定义资源仅用于数值累加准入控制
- 例如节点有
{"inference_license": 1.0},两个任务分别请求0.5可同时运行,各请求1.0则只能串行
- 例如节点有
- 注:不推荐将自定义资源用于节点属性筛选(如 GPU 型号),因其无法表达
不等于或包含等逻辑
标签选择器(属性型调度,Ray 2.49+)
官方推荐使用标签选择器处理非数值型调度约束 ,其机制独立于资源配额:
节点打标签 :
1
ray start --labels='{"region": "us-west", "instance_family": "spot"}'
任务声明约束 :在
.options()中传入NodeLabelSchedulingStrategy,支持In、NotIn、Exists操作符1
2
3
4
5
6from ray.util.scheduling_strategies import NodeLabelSchedulingStrategy
strategy = NodeLabelSchedulingStrategy(
label_selector={"region": "us-west", "instance_family": "spot"}
)
task_ref = my_function.options(scheduling_strategy=strategy).remote()
调度器在可行性检查阶段直接剔除标签不匹配的节点,再执行后续排序与随机选取
放置组(原子性跨节点资源预留)
- 当任务需要 同时成功预留跨多个节点的多份资源(如分布式训练:1个PS + 4个Worker),默认调度可能因部分资源预留成功、部分失败而陷入死锁
- 机制 :调用
placement_group()一次性向 GCS 声明资源包(Bundles)列表- GCS 执行原子调度,要么全部预留成功,要么全部回滚
- 绑定 :后续 Task/Actor 通过
PlacementGroupSchedulingStrategy绑定到特定 Bundle 上执行
核心配置建议
- Head 节点的 CPU :
- 建议设置
--num-cpus=0 - Head 节点仅负责任务调度和 GCS 状态维护,若允许用户任务占用其逻辑 CPU,可能导致系统组件响应延迟
- 直接使用 ray start 启动时,默认值为物理 CPU 核心数
- 建议设置
- Actor 的 CPU :
- 建议显式设置
num_cpus(如num_cpus=1) - 默认值为 0 会导致单节点可无限创建 Actor,耗尽物理内存
- 建议显式设置
- 内存资源 :
- 务必为任务设置
memory参数 - 若不设置,Ray 无法进行内存准入控制,高并发下极易触发节点 OOM 导致进程崩溃
- 务必为任务设置
- 浮点数精度 :
- 调度器进行浮点数累加
- 若使用
0.1等非精确二进制小数,多次加减后可能导致剩余容量与请求量的误差大于阈值,造成调度误判- 建议使用
0.5或0.25
- 建议使用
- 资源与属性分离 :
- 节点型号、区域等属性型约束请使用标签选择器 ,不要使用自定义资源(自定义资源仅保留给数值型并发控制)
补充:指定任务执行 Worker 的方法
- Ray 中有多种方法可实现
方法一:使用节点亲和性策略 (NodeAffinitySchedulingStrategy)
最直接的方法,通过节点的唯一 ID 将任务强制绑定 到特定节点
在提交任务或创建 Actor 时,通过
scheduling_strategy参数指定NodeAffinitySchedulingStrategy,并将node_id设为刚才获取的 ID1
2
3
4
5
6
7
8
9
10
11
12
13
14from ray.util.scheduling_strategies import NodeAffinitySchedulingStrategy
# 假设这是你想要运行任务的目标节点ID
target_node_id = "YOUR_TARGET_NODE_ID"
def my_task():
return "This task runs on a specific node."
# 将任务调度到目标节点
# soft=False 表示强制调度,如果目标节点不可用,任务将失败;如果设为 `True`,当目标节点不可用时,任务可以被调度到其他节点
result = my_task.options(
scheduling_strategy=NodeAffinitySchedulingStrategy(node_id=target_node_id, soft=False)
).remote()PS:查看集群中有哪些节点:
1
2
3
4
5
6
7
8import ray
ray.init()
# 获取集群中所有节点的信息
nodes = ray.nodes()
for node in nodes:
# NodeID 是节点的唯一标识符
print(f"Node ID: {node['NodeID']}, Resources: {node['Resources']}")
方法二:使用标签选择器 (Label Selector) 【推荐】
标签选择器是 Ray 2.49 版本后引入的更灵活、更强大的方式
标签选择器通过为节点打上标签,然后在任务中声明需要匹配的标签来选择节点
1)为节点打标签 :在启动 Ray 节点时,通过
--labels参数为其添加标签1
2# 启动一个 Worker 节点,并打上标签
ray start --address='YOUR_HEAD_NODE_ADDRESS' --labels='{"gpu_type": "A100", "region": "us-west"}'2)提交任务 :在任务中通过
label_selector参数指定标签选择条件1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
def my_task():
return "This task runs on a node with specific labels."
# 将任务调度到拥有 "gpu_type=A100" 标签的节点上
result = my_task.options(
label_selector={"gpu_type": "A100"}
).remote()
# 定义一个标签选择器,要求节点的 "instance-type" 标签不在 ["spot"] 中
from ray.util.scheduling_strategies import NodeLabelSchedulingStrategy
strategy = NodeLabelSchedulingStrategy(
label_selector={
"instance-type": {
"operator": "NotIn",
"values": ["spot"]
}
}
)
result = my_task.options(scheduling_strategy=strategy).remote()- 标签选择器支持
In、NotIn等复杂操作,可以实现 “运行在非 Head 节点” 或 “运行在region为us-west或us-east的节点” 等精细控制
- 标签选择器支持
方法三:使用放置组 (Placement Group)
放置组主要用于跨节点、多资源 的原子性预留,但也可以用来将任务“引导”到特定节点
1)创建放置组 :定义资源包(Bundle),并通过
strategy策略(如STRICT_PACK)控制其分布1
2
3
4
5
6
7from ray.util.placement_group import placement_group, PlacementGroupSchedulingStrategy
from ray.util.scheduling_strategies import PlacementGroupSchedulingStrategy
# 创建一个放置组,包含一个需要 1 个 CPU 资源的 Bundle
# STRICT_PACK 策略会尝试将所有 Bundle 放在同一个节点上
pg = placement_group([{"CPU": 1}], strategy="STRICT_PACK")
ray.get(pg.ready()) # 阻塞当前程序,直到指定的放置组在集群中成功创建并完成资源预留,若集群资源不足,该调用会永久阻塞,建议设置超时(如 pg.ready(timeout_seconds=30))2)获取 Bundle 所在节点 :通过
placement_group的 API 可以获取其 Bundle 被调度到的节点 ID3)绑定任务 :将任务绑定到该放置组,从而让它运行在放置组所在的节点上
1
2
3
4
5
6
7
8
def my_task():
return "This task runs where the placement group is."
# 将任务调度到该放置组所在的节点
result = my_task.options(
scheduling_strategy=PlacementGroupSchedulingStrategy(placement_group=pg)
).remote()
方法四:使用自定义资源 (Custom Resource) 【旧方案,不推荐】
这是一种旧方案,通过伪造“资源”来实现节点选择。例如,为特定节点设置一个
{"node-type": 1}的资源,然后在任务中请求该资源示例:
启动节点
1
2# 在目标节点启动时声明自定义资源
ray start --address='HEAD_ADDRESS' --resources='{"node-type": 1}'提交任务
1
2
3
def my_task():
return "This task runs on a node with the custom resource."
不推荐原因这种方式的原因是:
- 这种方法混淆了“资源(Resource)”和“约束(Constraint)”的概念,且只能表达“等于”的逻辑,缺乏灵活性
- 官方建议使用标签选择器作为替代方案