分布式系统(31):Ray 的 Future、对象血缘与重试边界
Ray 把远程函数的返回值立即表示成 ObjectRef,它在概念上类似论文所称的 future。调用者可以继续把这个引用交给下游任务,不必先等待真实值,于是普通 Python 函数调用扩展成执行期间不断生长的分布式 DAG。代价也随之出现:引用、对象值、生产任务和外部副作用有不同的故障命运。
RDD 以数据集转换为中心;Ray 的基本节点更细,可以是一项任意远程任务及其 future。两者都能沿依赖重算,但 Ray 当前产品的对象 ownership、task retry 和 actor 规则不能直接从 RDD 类比推出。
ObjectRef 是未来值,不是结果副本
f.remote() 提交一个 task 并返回 ObjectRef。下游 task 接收引用时,Ray 记录依赖;参数值可用后,下游才具备运行条件。ray.get 把等待显式带回调用者。
flowchart LR
A[f.remote → ObjectRef A] --> C[h.remote A B]
B[g.remote → ObjectRef B] --> C
C --> R[ObjectRef C]
R --> G[ray.get]
引用不是对象字节的副本。当前 Object Fault Tolerance 文档把对象数据放在 object store,把位置等元数据放在对象 owner。owner 通常是创建原始 ObjectRef 的 worker;实际计算值的 worker 可以是另一个进程。小对象还可能内联进进程内存或 RPC,因此这里讨论的是 Ray 的对象语义,不假定每个返回值都落入 Plasma。判断故障时必须问三件事:引用由谁拥有、值有哪些副本、值能否由生产 task 重建。
ObjectRef 作为顶层 task 参数时,Ray 会在任务执行前自动取出其值;嵌套在容器中的引用不会自动解引用。这一区别影响 DAG 是否真的建立了数据依赖,也避免把一个普通 Python 容器误认为运行时可追踪的依赖边。
DAG 同时表达依赖与可重建路径
Ray 论文把系统设计成面向动态 task graph 的分布式执行框架。task 的输出成为 future,下游引用把边连起来。只要生产 task 的描述及传递依赖仍可用,对象值丢失后就有一条重建路径。
flowchart LR
X[input-a=2] --> S[sum]
Y[input-b=3] --> S
S --> O1[sum-ref=5 @ node-a]
O1 --> M[scale]
M --> O2[scaled-ref=50 @ node-b]
本地模型先删除 node-b 上的 scaled-ref。sum-ref 尚存时,只需再执行一次 scale。再删除两个节点上的中间值时,请求 scaled-ref 会先重放 sum,再重放 scale。这展示的是依赖闭包,不是 Ray 真实调度轨迹。
对象重建还有明确限制。当前官方文档指出:ray.put 创建的对象没有可重放生产 task,不能用这套 lineage 重建;对象的 owner 必须仍存活;task及传递依赖需可重建;重试不能超过配置上限。actor task 结果默认也不可重建,除非显式设置 actor 的 task retry。
对象丢失与 owner 丢失不是同一故障
对象值所在节点丢失时,Ray 先尝试找同一对象的其他副本;没有副本时才重执行生产 task。owner 进程死亡则不同:当前文档明确说不支持 owner failure recovery,后续获取可能得到 OwnerDiedError,剩余值副本还会被清理以避免泄漏。
flowchart TD
L[ObjectRef取值失败] --> V{值有其他副本?}
V -->|有| Copy[读取副本]
V -->|无| Owner{owner仍活着?}
Owner -->|否| OD[OwnerDiedError]
Owner -->|是| Lineage{task lineage可重放?}
Lineage -->|是| Replay[递归重建]
Lineage -->|否| Lost[ObjectLostError类失败]
因此,“对象存储支持恢复”不能简写成“对象永不丢”。副本、lineage、owner 和 retry budget 共同决定恢复是否可行。
task failure 与应用异常采用不同策略
worker 进程或节点意外死亡时,当前 Ray 文档说明普通 task 默认最多重试 3 次,可用 max_retries 调整;-1 表示无限重试,0 禁止重试。函数抛出的应用异常默认不重试,需要用 retry_exceptions 显式选择。
这是合理的故障分类。worker 消失通常没有一份可信结果;业务异常可能是确定输入触发的稳定失败,盲目重试只会重复成本。无论哪类重试,客户端看到成功只说明某个 attempt 最终交付了结果,并不说明函数体只执行过一次。
sequenceDiagram
participant D as Driver
participant W1 as Worker attempt 1
participant X as External API
participant W2 as Worker attempt 2
D->>W1: task(input)
W1->>X: effect succeeds
Note over W1: worker dies before object is accepted
D->>W2: retry same task
W2->>X: effect succeeds again
W2-->>D: ObjectRef becomes ready
lineage reconstruction 也会重执行生产 task,因此同样触发函数体。重建正确性依赖 task 确定且幂等;Ray 不会替应用检查这两个性质。
调度决定放哪里,不决定业务语义
当前 Ray 调度先满足 task 声明的资源硬约束。未指定其他策略时,task 还会优先考虑大参数对象的数据局部性;显式选择其他 scheduling strategy 后,不再应用这项局部性偏好。placement group、node affinity、spread 等机制表达不同放置要求。
这些机制回答“哪个节点适合运行”。它们不回答输出是否已经外部提交,也不把 GPU、标签或 placement group 变成容错协议。软 node affinity 在目标节点死亡后可以转移;硬 affinity 的目标节点不存在或不可行时,task 会以不可调度错误结束。目标节点仍存活但资源暂时不可用,task 才会等待资源。选择必须与活性目标一致。
本地实验:删除对象与失败重试
实验位于 examples/distributed-systems/ray31/check.py。运行环境没有安装 Ray,因此这里只用确定性状态机模拟 input → sum → scale DAG。
1 | |
Python 3.12.3 的正式结果中,只丢叶对象时 sum 执行 1 次、scale 执行 2 次,重建结果从故障的 node-b 转移到 node-a;连传递依赖一起丢失时两者都执行 2 次,并在仍存活的 node-c 生成结果。模拟 task 在外部效果完成后失败,重试使效果计数变成 2。没有生产 task 的 put-ref 丢失后返回“no reconstructable producer”。
完整输出见结构化观察,运行与来源见实验证据,未验证边界见验证说明。
副作用需要独立提交边界
若 task 只从输入计算输出对象,重复执行通常可由确定性和幂等性吸收。若 task 向数据库递增、发送邮件或调用支付 API,Ray 无法用删除一个 object 或丢弃一个 attempt 来撤销已经发生的效果。
可选方案包括:用 task 输入和逻辑操作号构造 sink 幂等键;先生成不可变输出对象,再由单独提交者发布;为 actor 显式配置重启,并由应用自行实现 checkpoint、恢复与去重;使用外部事务资源。actor 默认不重启,重启时也只是重跑构造器,不会自动还原业务状态。每种方案只在自己的参与者集合内成立。
安全性、活性与恢复判断
纯函数 DAG 的安全直觉与第 30 篇相似:根输入确定,task 对相同参数产生相同输出,则拓扑顺序归纳可得重建值与原值相同。这里还需要 owner 与task描述可用、重试次数未耗尽。外部副作用不在这个归纳里。
活性要求存在满足资源约束的健康节点、依赖对象最终可获得或可重建、owner与控制面仍可服务,并且失败不会无限持续。无限 max_retries 只允许继续尝试,不保证某次尝试最终成功。
| 现象 | 应检查 | 不应直接推出 |
|---|---|---|
| ObjectRef ready | 计算已终止,值或错误已可观察 | 函数体只执行一次 |
| 对象值丢失 | 副本、owner、lineage、retry budget | 任意对象都能恢复 |
| task retry成功 | 某次attempt产生结果 | 外部副作用恰好一次 |
| task等待资源 | resource与scheduling strategy | 系统发生共识阻塞 |
| actor结果丢失 | actor restart/checkpoint/task retry配置 | 普通task规则自动适用 |
两个推演练习
只丢叶对象。 A=f(x)、B=g(A),A仍在另一节点,B唯一副本丢失。恢复B需要重跑哪些task?若A也丢失呢?
前者只需重跑 g;后者需先重建A再执行 g,前提是x、owner、两个task描述和重试预算仍有效。
副作用后崩溃。 task把订单状态推进后worker退出,Ray没有收到返回对象。把 max_retries 提高能否避免重复推进?
不能。它提高了得到结果的机会,也扩大了重复执行可能。sink必须用稳定操作标识去重,或把读取、写入和完成记录纳入同一个事务边界。
工程结论
Ray 的 ObjectRef 把数据依赖直接暴露给运行时;对象存储承载值,owner保存关键元数据,lineage提供重建配方,调度器选择可行节点。四者组合才能解释一次失败后的行为。
下一篇进入流处理。Ray DAG 的完成条件围绕future,流处理还要持续推进事件时间、水位线和有状态算子的checkpoint,失败恢复与外部sink原子性会更加突出。
参考资料
- Moritz et al., 2018, Ray: A Distributed Framework for Emerging AI Applications:§2–§4;描述论文时期的动态任务图与两级调度架构。
- Cheng et al., 2021, Ownership: A Distributed Futures System for Fine-Grained Tasks:owner-based distributed futures 与对象恢复边界。
- Ray Tasks:remote task、ObjectRef、依赖、调度与重试,访问 2026-09-26。
- Ray Task Fault Tolerance:系统失败、
max_retries与retry_exceptions,访问 2026-09-26。 - Ray Object Fault Tolerance:ownership、lineage reconstruction及限制,访问 2026-09-26。
- Ray Scheduling:资源、策略与数据局部性,访问 2026-09-26。
- Ray Actor Fault Tolerance:actor重启、task重试与checkpoint边界,访问 2026-09-26。
- Ray Task Lifecycle:基于Ray 2.48的task提交、参数传递、返回值与worker lease路径,访问2026-09-26。
