用资源约束限制 Ray 并发任务数量:基于 memory 与 num_cpus 的 OOM 防护实践
【免费下载链接】rayRay is an AI compute engine. Ray consists of a core distributed runtime and a set of AI Libraries for accelerating ML workloads.项目地址: https://gitcode.com/gh_mirrors/ra/ray
导读
在 Ray 中,任务的并行度由调度器根据资源可用性自动决定。但当任务需要加载大量数据到堆内存时,默认的并发策略可能让节点内存被迅速耗尽并触发 OOM。本文介绍 Ray 官方核心模式「使用资源(Resources)限制并发运行任务数」,通过调整任务的memory、num_cpus等资源需求,主动控制单节点上同时运行的任务与 Actor 数量,从根本上预防内存过载。读完本文,你将掌握资源需求与调度准入的底层机制、无限制与有限制两种写法的对比,以及 Actor、逻辑资源、自定义资源等边界场景的处理方法。
模式概述:为什么要限制并发运行的任务数
在 Ray 中,任务(Task)和Actor的并发度并非无约束的,调度器会依据每个节点可用的资源数量来决定能同时运行多少个任务。默认情况下:
- Ray 的每个任务默认请求1 个 CPU资源;
- Ray 的每个 Actor 默认请求0 个 CPU资源(调度时 1 个、运行时 0 个,详见后文)。
因此,调度器会把任务的并发数限制在可用 CPU 数量以内,而 Actor 的并发数则几乎是无限的。对于单个任务使用超过 1 个 CPU(例如通过多线程)的场景,并发任务之间可能因竞争而出现性能下降,但一般不会引发严重问题。
真正危险的是内存:如果任务或 Actor 使用的内存超过了它应有的份额,多个并发实例叠加起来就可能让节点过载,引发 OOM(内存耗尽)。针对这种情况,本模式给出的解决方案非常直接——提高每个任务或 Actor 请求的资源量,从而减少每个节点上并发运行的任务或 Actor 数量。
这一做法之所以有效,是因为 Ray 保证:同一节点上所有并发运行的任务与 Actor 的资源需求之和,不会超过该节点的总资源。这是 Ray 调度器在准入控制(admission control)环节强制保证的不变量,也是本模式的原理基石。
注意(Actor 任务):对于 Actor 上执行的任务(actor tasks),并发运行的 Actor 数量会直接限制可同时运行的 actor task 数量。也就是说,想限制 actor task 的并发度,先要限制 Actor 本身的并发度。
典型使用场景:数据处理负载的 OOM 防护
本模式最典型的应用场景如下:
你有一个数据处理工作负载,它使用 Ray 的远程函数(remote functions) 独立地处理每一个输入文件。由于每个任务都需要把输入数据加载进堆内存并完成处理,同时运行的任务过多时很容易导致 OOM。
在这个场景中,每个任务都需要大量堆内存,而 CPU 用量可能并不高。此时就可以用memory资源来限制并发运行的任务数(使用num_cpus等其他资源也能达到同样的目的)。
需要特别强调的是:与num_cpus类似,memory资源需求是"逻辑"的。Ray 不会强制监控每个任务的实际物理内存占用,也不会在任务实际内存超过声明值时干预或杀死任务。也就是说,资源需求只影响调度准入,真正的内存自律需要由任务自身保证——这是一个在使用本模式前必须理解的前提。
代码示例:无限制 vs 有限制
完整的可运行示例位于仓库中的 limit_running_tasks.py,下面对照展开。
无限制版本:16 个并发任务压垮 16G 内存
import ray # 假设该 Ray 节点有 16 个 CPU 和 16G 内存。 ray.init() @ray.remote def process(file): # 实际工作为读取文件并处理数据。 # 假设每个任务需要占用 2G 内存。 pass NUM_FILES = 1000 result_refs = [] for i in range(NUM_FILES): # 默认情况下,process 任务使用 1 个 CPU 资源,不请求其他资源。 # 这意味着 16 个任务可以并发运行, # 而 16 个任务总共需要 32G 内存(远超节点 16G),必然 OOM。 result_refs.append(process.remote(f"{i}.csv")) ray.get(result_refs)在没有限制的版本中,process任务只声明了默认的 1 CPU 资源,没有声明任何内存需求。节点有 16 个 CPU,所以调度器允许 16 个任务并发运行;而每个任务实际需要 2G 内存,16 个并发任务共需要 32G,远超节点 16G 的内存容量,最终会触发 OOM。
有限制版本:通过 memory 资源把并发数压到 8
result_refs = [] for i in range(NUM_FILES): # 现在每个任务声明使用 2G 内存资源, # 并发运行的任务数被限制为 8 个(16G / 2G)。 # 在这种场景下,把 num_cpus 设为 2 也能达到同样的效果。 result_refs.append( process.options(memory=2 * 1024 * 1024 * 1024).remote(f"{i}.csv") ) ray.get(result_refs)关键改动只有一处:通过process.options(memory=2 * 1024 * 1024 * 1024)为每个任务显式声明2G 内存资源需求(注意单位是字节,2 * 1024 * 1024 * 1024即 2GiB)。节点总内存为 16G,因此调度器最多允许16G / 2G = 8个任务并发运行,8 个任务总共最多消耗 16G 内存,恰好不会超过节点容量,OOM 问题随之消除。
代码注释中还给出了一个等价写法:把num_cpus设为 2 也能达到同样的并发限制效果(16 个 CPU / 2 CPU = 8 个并发任务),因为调度器只关心"资源需求总和不超过节点资源"这一约束,而不在乎具体使用哪种资源。
原理纵深:Ray 资源需求与调度准入机制
要深入理解本模式,需要掌握 Ray 资源系统的基本模型,详见文档 Resources。
资源是"逻辑"的,而非物理的
Ray 中的资源是一个键值对:键是资源名称,值是浮点数数量。CPU、GPU、内存是 Ray 原生支持的预定义资源,除此之外还支持自定义资源。
关键的语义是:Ray 资源是逻辑资源,不必与物理资源一一对应。例如可以通过ray start --head --num-cpus=0在物理上拥有 8 个 CPU 的机器上启动一个拥有 0 个逻辑 CPU 的 head 节点(这样做主要是让调度器不在 head 节点上调度任务,把资源留给 Ray 系统进程)。
逻辑资源的语义带来了几个重要推论:
- 资源需求不约束实际物理资源使用:Ray 不会阻止一个
num_cpus=1的任务启动多线程并使用多个物理 CPU;同样,声明memory=2G的任务若实际吃掉 4G 内存,Ray 也不会干预。确保任务实际用量不超过声明值,是开发者自己的责任。 - Ray 不做 CPU 隔离:不会为某个任务预留物理 CPU 或做亲和性绑定,线程调度交给操作系统(必要时可用
sched_setaffinity等系统 API 自行绑定)。 - Ray 提供 GPU 隔离:通过自动设置
CUDA_VISIBLE_DEVICES环境变量以"可见设备"的形式隔离 GPU(详见 Accelerators)。
此外,若为任务/Actor 设置了num_cpus,Ray 会为其设置OMP_NUM_THREADS=<num_cpus>环境变量;未指定时则设为 1,以避免多 worker 场景下 OpenMP 线程争抢。这一点在使用多线程线性代数库(numpy、PyTorch、TensorFlow)时值得留意。
节点资源的默认配置规则
在不做任何显式配置时,Ray 会按以下规则自动确定每个节点的逻辑资源量:
- 逻辑 CPU 数(
num_cpus):等于机器/容器的 CPU 数; - 逻辑 GPU 数(
num_gpus):等于机器/容器的 GPU 数; - 内存(
memory):等于 Ray 运行时启动时"可用内存"的70%; - 对象存储内存(
object_store_memory):等于"可用内存"的 30%(注意:对象存储内存不是逻辑资源,不能用于调度)。
因此,本文示例中"节点有 16G 内存"对应的实际逻辑memory资源量约为 11.2G(16G × 70%),读者在自己环境中复现时,应结合实际节点内存与配置计算并发上限。注意 Ray不允许在节点启动后动态更新资源容量。
调度保证:并发资源需求之和不超过节点总量
文档 Resources 中明确给出了任务与 Actor 资源需求对调度并发度的影响:
同一节点上所有并发执行的任务与 Actor 的逻辑资源需求之和,不能超过该节点的总逻辑资源。
这正是本模式的全部依据:把每个任务的内存需求声明得越大,能被同时调度到该节点上的任务就越少。文档还专门把"利用这一性质限制并发运行任务、规避 OOM"的做法,链接到了本文讲解的模式(即core-patterns-limit-running-tasks),说明这是官方推荐的标准实践。
用 options() 与 ray.remote() 声明资源需求
除了示例中的process.options(memory=...),还有多种声明资源需求的方式:
- 在装饰器中声明:
@ray.remote(num_cpus=2, memory=2 * 1024 * 1024 * 1024); - 在调用时声明:
process.options(num_cpus=2, memory=...); - 使用自定义资源:
process.options(resources={"special_hardware": 1})。
Java 与 C++ 也支持等价写法,例如 C++ 中为ray::Task(MyFunction).SetResource("CPU", 1.0)、ray::Actor(CreateCounter).SetResource("GPU", 1.0)。
分数资源需求
Ray 支持分数资源需求。例如任务/ Actor 是 IO 密集型、CPU 占用很低时,可以指定num_cpus=0.5甚至num_cpus=0。分数资源需求的精度为 0.0001,应避免使用超出该精度的浮点数。另需注意:GPU、TPU、neuron_cores 资源需求若大于 1,必须为整数,例如num_gpus=1.5是非法值。
Actor 场景与默认值陷阱
文档 Actors 与 Resources 中明确指出:
默认情况下,Ray 任务使用 1 个逻辑 CPU 资源;Actor 在调度时使用 1 个逻辑 CPU,在运行时使用 0 个逻辑 CPU。
这意味着,默认情况下 Actor 无法被调度到 0 CPU 节点上,但在任何非 0 CPU 节点上都可以无限量地运行(这一默认行为是历史原因造成的)。因此:
- 想限制 Actor 的并发数,必须显式设置
num_cpus(例如@ray.remote(num_cpus=1)),否则并发 Actor 数量不受 CPU 约束,也就无法借助本模式控流; - 若显式声明了资源需求,这些资源在调度时和运行时都会被要求满足;
- 对于 actor task,限制 Actor 本身的并发数即是限制其任务的并发数(本模式开头 note 的内容)。
与其他限流模式的区分:limit-pending-tasks
Ray 的核心模式索引中还有一个容易混淆的模式「使用 ray.wait 限制 pending 任务数」。两者的分工如下:
- 本文模式(limit running tasks):通过调整每个任务的资源需求,控制有多少任务可以同时运行,属于调度层面的资源准入控制,是本模式推荐的"修改并发运行数"的标准方法;
- limit-pending-tasks 模式:通过
ray.wait()施加背压(backpressure),控制有多少任务处于在途(in-flight)/ 排队(pending)状态,防止无限提交任务导致 pending 队列无限增长。它主要用于限制"同时飞行中的任务数",虽然也可用于限制并发运行数,但不推荐,因为会影响调度性能。
两者的适用场景也不同:如果一次性提交有限数量的任务,每个任务在队列中仅占用少量内存做簿记,通常不会出问题;真正危险的是源源不断提交任务的无限流。当面对无限任务流时,优先考虑ray.wait()背压;而当目标是精确控制并发运行数时,则应回归本文的资源需求方案。
小结与最佳实践清单
总结本模式的实践要点:
- 默认并发:任务默认 1 CPU,Actor 默认运行时 0 CPU;CPU 限制了任务并发,但 Actor 并发近乎无限。
- OOM 根因:任务/Actor 实际内存用量超过其资源份额,多实例叠加导致节点内存耗尽。
- 核心手段:通过
task.options(memory=...)、num_cpus或自定义资源提高单任务资源声明,利用"节点上并发资源需求之和 ≤ 节点总资源"的调度不变量压降并发数。 - 逻辑资源语义:
memory与num_cpus都是逻辑资源,Ray 不强制监控实际物理用量,任务自身需自律。 - 节点资源默认值:内存逻辑资源默认是可用内存的 70%,复现示例时需按实际节点配置计算并发上限。
- Actor 需显式设
num_cpus,否则并发不受控;actor task 的并发由 Actor 并发决定。 - 区分两种限流:控制"运行中"任务数用本模式(资源需求),控制"排队中"任务数用 ray.wait 背压模式。
通过为任务声明合理的内存(或 CPU)资源需求,你可以在不改动任何业务逻辑的前提下,用一行options()调用为数据密集型负载建立可靠的 OOM 防线。
【免费下载链接】rayRay is an AI compute engine. Ray consists of a core distributed runtime and a set of AI Libraries for accelerating ML workloads.项目地址: https://gitcode.com/gh_mirrors/ra/ray
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考