Ray 把远程函数的返回值立即表示成 ObjectRef,它在概念上类似论文所称的 future。调用者可以继续把这个引用交给下游任务,不必先等待真实值,于是普通 Python 函数调用扩展成执行期间不断生长的分布式 DAG。代价也随之出现:引用、对象值、生产任务和外部副作用有不同的故障命运。

分布式系统(30):Spark/RDD 的血缘、物化与故障重算

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
2
3
4
5
mkdir -p examples/distributed-systems/.build/ray31/tmp
export TMPDIR="$PWD/examples/distributed-systems/.build/ray31/tmp"
export TMP="$TMPDIR" TEMP="$TMPDIR" PYTHONDONTWRITEBYTECODE=1
python3 -B examples/distributed-systems/ray31/check.py \
--output examples/distributed-systems/.build/ray31/observations.json

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原子性会更加突出。

参考资料