Ray Task 机制详细学习笔记
Ray 中所有用
ray.remote装饰的函数都会变成异步 task,其返回值统一为 object reference;围绕ray.get、ray.wait、num_returns、generator、dynamic、cancel、异常处理与资源释放,形成了一套任务调度与结果获取机制。
一、核心要点速览
- Ray 是一个常用于强化学习领域的框架,本次按 Ray 文档梳理其中最核心的 task 相关概念。
- 加
ray.remote装饰器的函数会变成可远程调用的 task;调用方式是在函数名后加.remote。 - task 可以异步执行,主程序不会因为 task 内部的耗时操作而阻塞。原文用
sleep 10秒的slow function演示,运行时主程序没有立即被卡住。 - task 的返回值一定是 object reference,不论返回的是单个值、list、tuple 还是 generator。
- 获取 object reference 中的具体值要用
ray.get;ray.get会阻塞。 - 多个 task 可以用
ray.wait批量等待,返回两个 list:已完成部分和剩余部分。 - 如果希望返回多个值,可以在函数上设置
num_returns,也可以在 options 中设置;如果不设置,通常返回一个 object reference,指向 list 或 tuple。 - generator 场景中,如果长度动态,可以设置
num_returns="dynamic";此时会得到 object ref generator,可以遍历其中的 object reference,再逐个ray.get取值。 - 使用 dynamic generator 时,原文建议把重试次数设为 0,即使用
max_retries=0,因为异常重试后的返回结果不可信。 - 调用
cancel后,相关 object reference 不能再被ray.get,否则会报错,例如object reference was cancelled。 - task 执行时会释放 CPU resource;
ray.get的核心 resource 位于 get block 内部。如果在 block 中继续占用资源,容易造成死锁。
二、Task 与 ray.remote
2.1 Task 的类比与底层理解
原文把 Ray 的 task 类比为传统异步调度中的任务:
- 如果按 C++ 的理解,可以把它理解成 future。
- 底层应该也是基于 event loop 的模式实现的。
- 按照文档说法,task 可以异步执行。
这里“异步”的直接表现是:调用 task 后不会立刻拿到具体结果,而是先得到一个 object reference;主程序可以继续执行,不阻塞等待。
2.2 slow function 示例
原文举了一个例子:
- 定义一个
slow function。 - 在函数上使用
ray.remote装饰器。 - 使用
ray.remote后,这个函数会变成可以被远程调用的函数。 - 调用方式是
slow.remote()。 - 如果
slow和调用方不在同一个节点上,Ray 会进行远程调用。 - 运行示例时,因为函数内部有
sleep 10秒,主程序不会阻塞等待;原文观察到运行时没有立刻输出结果,但主程序并没有被卡住。
2.3 资源分配
原文提到,可以在 remote 上分配资源。但作者当时没有 GPU 机器,所以没有实际演示资源分配。原文也强调资源分配是一个比较重要的领域,计划在后面的 resource 部分再展开。本次 task 部分只先记录这一点。
三、Object Reference 与 ray.get
3.1 Object Reference 是什么
Object reference 是 remote 调用之后返回的对象。原文通过打印观察到了以下现象:
- 如果加
print,会看到 object reference。 - 在示例中打印了四次。
- 每个打印出的地址或数值结构看起来不太一样。
- 原文表示具体这个结构是什么暂时不太清楚。
关键结论:remote 调用不会直接返回函数执行结果,而是返回 object reference。
3.2 如何获取值:ray.get
获取 object reference 中具体值的方法是 ray.get:
- 写法是
ray.get(object_ref)。 ray.get是阻塞的。- 原文实验中发现,如果直接
ray.get一个耗时 10 秒的 task,就会等待 10 秒。 - 为了避免一次等 10 秒,原文改成
sleep 1秒,然后观察结果,发现它会每隔一秒钟打印一个结果。
这说明:多个 task 可以并行提交,但逐个 ray.get 时仍然会按阻塞方式取结果。
四、ray.wait 与批量任务调度
4.1 为什么需要 ray.wait
如果有多个结果,直接用 ray.get 会阻塞住,直到结果完成。原文希望实现的效果是:
- 每个 function 相当于一个 task。
- 把所有 task 丢到一个调度器里,让 Ray 帮忙执行。
for循环不是串行执行,而是并行执行。
这时可以使用 ray.wait。
4.2 ray.wait 的写法与返回值
原文描述的写法大致是:
- 先构造一批 task,例如把每个
i传给f.remote(i)。 - 然后用
ray.wait等待这些 task。
ray.wait 的返回结果:
- 返回两个 list。
- 第一个 list 包含已经完成的部分。
- 第二个 list 包含剩余未完成的部分。
- 默认
return value是 1,也就是默认等待完成 1 个任务就返回。 - 如果希望所有结果都返回,可以设置
num_returns=len(tasks)。 - 设置为全部等待后,实验中发现未完成 list 为空。
4.3 ray.wait 与 ray.get 的对比
| 操作 | 行为 | 返回 | 适用场景 |
|---|---|---|---|
ray.get | 阻塞等待结果 | 具体值 | 获取单个或多个结果 |
ray.wait | 等待部分任务完成 | 两个 list:已完成、剩余 | 批量任务,不想一次全部阻塞 |
ray.wait 默认 | 默认等待 1 个完成 | 已完成 list、剩余 list | 希望先处理已完成任务 |
ray.wait(num_returns=len(tasks)) | 等待全部完成 | 剩余 list 为空 | 等所有 task 完成 |
原文还提到一种阻塞做法:阻塞住直到四个 task finish。这和 ray.wait 类似,只是 ray.wait 更灵活,可以设置等待完成的数量或长度。
五、多返回值与 num_returns
5.1 multiple returns 的行为
原文讨论 multiple returns 时说:
- 如果是 multiple,会建立一个 multiple object reference。
- 如果有一个 return value,会给它创建一个 return value 的 reference。
- 如果设置
num_returns,会区分出 multiple reference。 - 如果不设置,可能会当成一个 tuple 去做。
5.2 设置方式
原文中提到的设置方式包括:
- 在函数上设置
num_returns。 - 在 options 中设置
return values。 - 通过
ray.get解包。 - 使用 generator 时设置
num_returns="dynamic"。
原文后来做了一个修正:
- 之前认为必须给
num_returns,但后来发现报错可能和多打了一个括号有关。 - 报错消息和
num_returns可能没有关系。 - 如果设置
num_returns=4,可以打印出四个 object ref。 - 如果不设置,不论返回的是 tuple 还是 list,都会返回一个 object reference。
- 这个 object reference 指向的是一个 list 或 tuple。
- 如果想把它解析出来,需要
ray.get解包。
5.3 解包过程
原文通过 g.remote 的例子说明:
g.remote可能返回一个 list 的 object reference。- 如果设置
num_returns=4,会返回四个结果。 - 如果不设置,返回的是一个 object reference,指向 list 或 tuple。
- 用
ray.get(g.remote())可以把这个 object reference 解包出来。 - 解包后得到的结果同样是一个 list。
- 因此,
ray.get在这里的作用就是一个解包过程。
六、Generator 与 dynamic
6.1 generator 示例
原文举了一个 generator 例子:
a、b、c分别是三个 reference。- 如果按某种方式直接返回,会报错。
- 如果给它一个 return reference,也依然报错。
- 原文据此推测:对于 generator 这种方式,Ray 在某些过程中知道 generator 的长度,应该不是一个动态过程。
- 如果 generator 长度无法获取,Ray 就很难判断出这种错误。
6.2 dynamic 设置
更常见的场景是不知道 generator 的具体返回值是多少,也就是动态长度。这时可以:
- 设置
num_returns="dynamic"。 - 如果不设置,可能会报错,例如
on hit error task ... through exception ...。 - 设置 dynamic 后,可以正常运行。
- 原文观察到打印了
return value 0,并返回了一个 object ref,即created object。 - 前面某些场景下反而打印了返回值,原文表示不太清楚 Ray 内部如何判定。
6.3 object ref generator 的遍历
使用 dynamic 时,得到的是一个 object ref generator:
- 它可以被用来迭代所有的 object reference。
- 可以把它当做一个 generator 去遍历。
- 遍历得到的是 object reference。
- 想获取 object reference 中的值,仍然要用
ray.get。 - 原文实验表明,这样可以正常拿到打印值。
整体逻辑是:
- dynamic 返回 object ref generator。
- 遍历 object ref generator,得到每个 object ref。
- 用
ray.get获取每个 object ref 中的值。
6.4 作为参数传入另一个 task
原文还提到:
- object ref generator 也可以作为一个 argument 传入到另一个 task 中。
- 这样做的根本原因是:如果通过
ray.get去获取一个 object ref generator,会阻塞在这里。 - 如果直接把 object ref generator 给到另一个 task,只需要在最外层获取具体 object ref。
- 这样可以避免中间的 block。
6.5 generator 的阻塞行为
原文强调了一个细节:
ray.get一个 generator 时,行为和 list 相同,都是阻塞式的。- 原文例子中,每个迭代过程中
sleep5 秒,实际等了 10 秒才获取到 generator。 - object generator 实际上是阻塞获取所有过程。
- 因此,它的行为更像是一个 list。
七、Cancel 机制
原文认为 cancel 是一个很有意思的问题:
- 如果深入了解过 TensorFlow 或一些 RPC 框架,会发现它们都有 cancel 机制。
- cancel 给了一个机会去取消之前做的一些操作。
- 在 Ray 中,任务会变成 task,task 会放到一些队列里面。
- 取消操作就是把 task 从队列里面移除掉。
- 对于 TensorFlow 这种计算图形式会更复杂:取消动作会一级一级向上影响;如果某个节点被 cancel,会一直影响到最上层的节点。
Ray 中的实验现象:
- 如果调用了
cancel,这个 object 就没办法get到了。 - 原文解释:
print应该不会有结果,或者应该会报错。 - 实验果然报错,报错内容与
object reference was cancelled相关。 - 原文还提到报错位置在 32 行,之前位置出错过。
八、嵌套 remote 与 option
8.1 在 remote function 中再调用 remote
原文介绍了一个更复杂的例子:
- 在一个 remote function 里面调用第二个 remote 过程。
- 例子中调用
g function。 g.remote返回的是一个 list 的 object reference。- 按之前说法,不论返回 list 还是其他东西,都会返回 object reference。
- 如果说明
num_returns,例如设置为 4,则返回四个结果。 - 因为底下的
for循环有四次,所以逻辑上可以返回四个 object reference。 - 如果不设置,则返回一个 object reference,指向 list 或 tuple。
ray.get(g.remote())得到的结果同样是一个 list。- 因此,
ray.get的作用就是解包。
8.2 option 传参
原文还发现可以用 option 传参:
- 可以不把参数写在
remote里面,而是写在 option 里面。 - 例如在 generator 调 remote 时给它一个 option。
- 实验表明,把参数写在 option 里面也是可以的。
- 前面遇到的一系列报错,后来发现可能是因为多打了一个括号,应该把括号去掉。
- 原文特别提醒:之前有一条报错消息和
num_returns没有关系,之前的结论可能也有一些问题。 - 通过 options 可以设置
return values=2,也可以设置 dynamic。
九、异常处理与 retry
原文最后讨论了异常处理:
- 前面运行了两个
for循环,执行完后又执行了一次报错异常。 - 默认让它执行四次,剩余的那些 task 会抛出一个异常。
- 使用 dynamic 时,可以获取到
ref 3。 - 例如
1、2是有值的,3获取到的是最后一个 object ref。 - 这个最后的 object ref 返回的是一个错误的值,也就是异常结果。
- 如果把这个 value 不捕获,会报一个 generator 的 pid 和 error,例如
unhandled load error。 - 如果使用 dynamic,则不会报任何错。
- 原文理解:报错应该是在做
ref 3时发生的错误。 - 如果使用 dynamic,应该设置重试次数。
- 原文提到这里符号好像没有,现在好像叫
max_retries。 - 文档这里可能有点问题。
- 要把重试次数设置成零,即
max_retries=0,因为如果出现异常还让它重试,返回的结果是不可信的。
十、资源释放与死锁
原文在 task 部分最后补充了资源释放问题:
- task 在被执行的时候,会释放 CPU resource。
- 当调用
ray.get时,核心 resource 实际上在 get block 里面。 - 如果在 block 里面执行一些运算操作,会给它 resource 去运行。
- 在
ray.get的地方,会把这个 resource 释放掉。 - 因为 Ray 会 block 住运行,如果此时资源还占着,容易造成死锁。
- 原文强调:虽然前面没有详细谈资源设置,但后面可能不会再专门说资源释放话题,所以这里先提一下。
十一、易错点与原文修正
- 报错来源可能不是
num_returns,而是多打了一个括号。原文在实验后修正了之前的判断。 - 如果不设置
num_returns,不论返回 tuple 还是 list,都会返回一个 object reference,指向 list 或 tuple。 - 如果设置
num_returns=4,可以打印出四个 object ref。 ray.get(g.remote())的行为是解包,返回 list。- generator 直接返回某些形式会报错;给 return reference 也可能报错。
- 使用 dynamic generator 时,如果不设置
max_retries=0,异常重试后的结果不可信。 - 调用
cancel后,object reference 不能再get,会报object reference was cancelled。 ray.get在 generator 上与 list 一样是阻塞的,行为更像 list。
十二、限制与待确认问题
以下问题在原文中没有完全展开或存在不确定:
- object reference 打印出的具体结构、地址或数值含义,原文表示暂时不太清楚。
- generator 的长度为什么能被获取,原文只是推测,没有给出明确机制。
- 报错与
num_returns的关系,原文先认为有关,后又修正为可能与括号错误有关,最终结论没有完全定论。 max_retries参数名与文档中出现的符号不一致,原文认为文档可能有问题。- 资源分配部分没有实际演示,因为作者没有 GPU 机器。
- 错误容忍机制只是提到,没有展开。
- TensorFlow 计算图的 cancel 影响只是类比说明,没有深入展开。
- dynamic 下为什么不报错,原文没有给出底层机制。
- object ref generator 的阻塞行为与 list 相同,但为什么“更像 list”没有进一步解释。
ray.wait中“设置等待长度”的表述,原文口播为“主材的长度”,具体所指可能需要结合文档确认。
十三、行动清单
- 给函数加
ray.remote装饰器,使其变成 task。 - 使用
函数名.remote()进行远程调用。 - 用
ray.get(object_ref)获取 task 返回值,并注意它会阻塞。 - 批量任务优先考虑
ray.wait,通过num_returns控制等待完成的数量。 - 需要多个返回值时,明确设置
num_returns或在 options 中设置return values。 - 用
ray.get解包 object reference 指向的 list 或 tuple。 - 使用 generator 时,若返回长度动态,设置
num_returns="dynamic"。 - 使用 dynamic generator 时,设置
max_retries=0,避免异常重试导致结果不可信。 - 异常场景中,注意通过后续 object ref 捕获异常,避免未处理错误。
- 调用
cancel后不要继续get对应 object reference。 - 注意
ray.get阻塞与资源释放之间的关系,避免死锁。 - 需要资源分配时在 remote 上指定,但本次原文没有展开具体操作。
十四、总结
Ray 中所有加了 ray.remote 装饰器的 function 都会变成 task。每个 task 的返回结果一定是一个 object reference,不论是 generator、list 还是 tuple,都是如此。
如果想返回多个值:
- 可以在函数上设置
num_returns。 - 也可以在 options 里设置。
- 还可以通过
ray.get解包,把 list 中所有 object 拿出来。 - generator 也是类似做法。
generator 需要特别注意:
ray.getgenerator 时,行为和 list 相同,都是阻塞式。- 异常处理时,异常会向后传递,可以通过后面的 ref 捕获异常。
- 对于 dynamic 方式,应该让
max_retries变成 0,不进行后续 retry 动作,因为原文认为 retry 返回的结果是错误的。