KubeRay 源码阅读:Operator 重启后如何继续工作
第三篇提到,一轮协调中创建资源成功,后面的状态更新仍可能失败。假设 demo-job 的 submitter Job 已经创建,Operator 却在写回 RayJob 状态前退出了。新 Operator 读到的仍是 Initializing,它会再提交一次计算吗?如果丢失的是 Ray Jobs API 的提交响应,答案又有什么不同?
本文继续使用第一篇的例子:RayJob 新建专用 RayCluster,通过 K8sJobMode 提交镜像中的程序,完成后等待 60 秒清理集群。讨论重试时,另设 spec.backoffLimit: 1,允许失败后再执行一次;submitter 使用默认命令。下文把一次完整的作业执行称为“执行轮次(attempt)”:它会经过多轮 Reconcile,可能跨过 Operator 重启,也可能包含提交程序自身的重试。
源码固定为 KubeRay 6bf05eb17a3e 与 Ray 3f785c0711b9。以下故障窗口均由源码推演,没有进行故障注入实验。五处代码均为连续真实节选,仅调整缩进。
1. 重启后还能读到哪些状态
先看只有 Operator 进程退出的情况,Kubernetes 控制面和 Ray 集群仍正常。旧进程的缓存、工作队列和函数执行现场随之丢失,RayJob、RayCluster 等 API 对象却还在,Ray 程序也可能一直在运行。
新 Operator 能接着做什么,取决于这些信息分别保存在哪里。图 1 把进程内状态与外部保存的状态分开。
| 保存位置 | 包含什么 | 对恢复的作用 |
|---|---|---|
Kubernetes 中的 spec |
入口命令、集群配置、重试与清理要求 | 说明用户希望执行什么 |
Kubernetes 中的 status |
管理阶段、submission ID、集群名、失败次数、时间与观测结果 | 保留已选定的身份和流程进度,部分观测可能滞后 |
| RayCluster、Job、Pod、Service 等对象 | 已创建的资源及各自状态 | 检查上次动作是否已经产生结果 |
| Ray 侧 | GCS 中的作业元数据、driver 和 actor 进程状态、对象数据 | 提供计算的实际进展,不能从 RayJob 状态重建全部内容 |
| Operator 内存 | 缓存、队列、局部变量 | 重启后重新建立 |
检查资源时,还要继续看它的状态。Pod 对象存在,容器可能仍未正常运行;RayJob status 中保存的 RayCluster 状态,也可能只是上一次观测的结果。
新进程启动监听并同步已有对象后,初始同步会把符合监听与筛选条件的对象交给事件处理器,再形成待协调请求。这就是第二篇通知链路在重启时的入口。RayJob 的 SetupWithManager 注册了 RayJob 及相关子资源;每轮 Reconcile 又按 namespace/name 读取 RayJob,再根据其状态选择分支。后续查询还会校正过时的观测值。
已有对象可以重新进入协调流程,所以恢复不依赖逐项还原旧队列。status 则不能随意清空,里面的身份字段决定了下一轮要查找哪个集群、哪个 Ray 作业。
2. 状态机记录一次执行走到了哪里
RayJob 有两个容易混淆的字段。jobDeploymentStatus 表示 KubeRay 管理流程的阶段;jobStatus 表示 Ray 入口程序的状态。创建 submitter Job 后,前者即可进入 Running,此时后者可能还没有变成 Ray 的 RUNNING。
图 2 展示专用集群、K8sJobMode 下的主要路径,并把重试画成新的执行轮次。
New 在 API 中用空字符串表示。Controller 先为它初始化身份并写入 Initializing,后续轮次才准备集群和 submitter。进入 Running 后,Controller 检查 submitter,再用 GetJobInfo 查询 Ray。沿正常终态判断路径,Ray 作业已经结束且 submitter 也已结束,才能收尾;submitter 失败或期限检查也能提前使管理流程失败。
这个版本支持 RayJob 级重试。状态分支推进到统一写回处时,checkBackoffLimitAndUpdateStatusIfNeeded 更新失败次数,并检查是否还有重试额度。对于 backoffLimit: 1,第一次可重试失败会写成 Retrying,第二次失败才耗尽额度。DeadlineExceeded 和 JobStatusCheckTimeoutExceeded 被明确排除在重试之外。Running 分支附近仍有“不支持 retries”的旧注释,应以这些执行分支为准。
进入 Retrying 后,Controller 删除旧 RayCluster 和 submitter Job,等到查询确认两个对象都不存在,再清空本轮集群名、submission ID 等状态,回到 New。这里检查的是对象删除完成;submitter 使用后台级联删除,即先删除上级 Job,再由垃圾回收器异步清理其 Pod,不能据此断言所有下级 Pod 都已退出。
这时再看两个名字相近的配置,就容易区分了:spec.submitterConfig.backoffLimit 控制 Kubernetes Job 重试提交程序,继续使用当前 Ray 作业身份;spec.backoffLimit 控制整次 RayJob 执行的重试,需要重新准备执行环境。
Operator 重启本身不会记作执行失败,也不会直接开启新的执行轮次。如果这段时间里 Ray 已经完成程序,接管后的 Controller 可以查到结果,再把状态补写到 Kubernetes。若读到的已是 Retrying,它则继续本轮清理;旧资源尚未删除时,不会先清空身份去创建下一套资源。
图中还列出几条主案例之外的分支。暂停先经过 Suspending 清理,再进入 Suspended,恢复时回到 New。InteractiveMode 把实际提交留给使用者,Waiting 是相应的等待阶段;配置校验失败则进入 ValidationFailed。正常终态映射中,Ray 的 STOPPED 可以对应管理状态 Complete,所以成功应核对 jobStatus: SUCCEEDED。使用已有集群的 clusterSelector 模式不允许配置正数的 RayJob 重试额度。
3. 先保存身份,再创建资源
同一个执行轮次里,每轮协调都要找到同一套资源。主案例的查找依据是 status.rayClusterName 中的 demo-job-abcde,以及 status.jobId 中的 submission ID;前者用于 Kubernetes 查询,后者用于 Ray Jobs 查询。
initRayJobStatusIfNeed 只在状态中的 submission ID 为空时选择它:
1 | |
同一函数也会选定 RayCluster 名称。自动生成的名称带随机后缀,随后与 Initializing 一起写回 status。New 分支不会接着执行 Initializing 分支,因此身份尚未写入成功时,这条路径还没有创建集群。
如果身份已经写入、响应却丢失了,后续重新读取仍能拿到原来的名字。如果写入没有成功,可以重新初始化,因为此时还没有创建后续资源。先保存名字,再按这个名字创建资源,使这两种情况都有可追踪的处理依据。
后续 getOrCreateRayClusterInstance 按 status.rayClusterName 查询:找到就复用,确认不存在才构造并创建;其他读取错误直接返回。已有集群模式下,找不到集群会报错,不会代建。
这也解释了另一种中断:RayCluster 创建成功,下一次状态写回却失败。因为集群名更早就已保存,新 Operator 仍用原名查找,无须凭本地日志猜测上次创建了什么。前提是 RayJob 的这份身份记录和对应资源仍在;人为清空 status 或删除、替换资源,会改变恢复所依据的事实。
重试开启新的执行轮次时,这些身份可以改变。自动生成的集群名和 submission ID 会重新生成;显式设置的 spec.jobId 则仍会被采用。submitter Job 一直使用 RayJob 的名字 demo-job,旧对象删除后新建对象会有新的 UID。UID 是 Kubernetes 为每个对象分配的唯一标识,所以名称相同也可能已经是另一个对象。稳定的是同一执行轮次内的查找依据,不能把某个名称当作跨重试、跨集群的永久执行凭证。
4. 创建成功,状态却没写回
回到开头的故障。集群已经就绪,Controller 在 Initializing 分支执行:
1 | |
最后一行只修改内存,状态写入 API 的动作在分支结束之后。创建 Job 与更新 RayJob 是两个请求,中间存在图 3 的故障窗口。
按以下顺序推演,时间间隔只作示意:
- API Server 保存了 submitter Job
demo-job,Kubernetes 已可据此启动提交程序。 - Operator 在更新 RayJob 前退出,或状态更新失败,RayJob 仍显示
Initializing。 - 新 Operator 再次处理该分支。
createK8sJobIfNeed查询同名 Job,读到已有对象便返回成功,然后继续推进状态。
若丢失的是创建响应,Operator 当时无法判断服务端是否成功,恢复仍靠查询同一名称。缓存暂时没看到新对象时,可能再次发出创建请求;同类型、同命名空间的同名对象不能同时创建,服务端会拒绝冲突,后续协调再读。
这条路径能在重复协调时复用已创建的资源,但依赖身份记录仍可读取、已有对象仍可识别。getOrCreate 中的查询和创建依然是两个操作,Kubernetes 写入、Ray 提交与用户程序向外部写结果也没有组成事务。因此,资源能被复用,并不足以承诺业务执行 exactly-once,也就是一次业务操作的效果恰好生效一次。
也不能把这个例子推广为“任何阶段缺少 submitter 都会补建”。在 Running 分支中,Controller 读取 submitter Job 的错误会返回,并没有调用前面的创建函数;Ray 中查不到作业时,K8sJobMode 也不会直接改走 HTTP 提交。恢复动作由当前阶段决定,排查时必须同时看保存的阶段和实际对象。
5. 提交响应丢失时,谁来防止重复
前面的故障发生在 Kubernetes 资源创建时。如果再往后一步,submitter 已经向 Ray 提交程序,却没有收到响应,随后提交 Pod 被重试,会发生什么?默认脚本会先查询当前 submission ID,查得到就跳过提交,继续跟踪日志。
BuildJobSubmitCommand 把查询放进 shell 条件中:
1 | |
这里的 jobStatusCommand 是 ray job status,提交命令则携带同一个 --submission-id。但 shell 只判断命令退出码:查询因网络故障而失败,也会进入提交分支。两个提交者还可能都先查到“不存在”。健康检查不能消除之后发生的故障,先查再提交也无法独自关闭并发窗口。
Ray 服务端还有一道检查。Jobs 服务用 supervisor 管理入口进程的启动和状态;JobManager.submit_job 在启动这个管理者之前,先尝试登记作业:
1 | |
底层把记录写入 GCS 的内部键值存储接口(internal KV),不覆盖已有键。同一记录仍在时,重复提交会被拒绝;它也不会把重复请求直接当作成功返回。因此,客户端预查询改善的是重试时的处理路径,服务端记录提供另一层防重依据。
这项检查只在对应记录仍然存在时有效。换了集群、换了 submission ID,或者原记录已经删除、丢失,原有防重依据就不再适用。记录写入与进程启动之间还隔着一步,因此查到记录也不能证明程序已经执行成功。使用自定义 submitter 命令时,还要确认是否保留了默认预查询逻辑。
例如,程序已经往外部数据库写入一批结果,随后失败,RayJob 重试启动了新的执行轮次。这次执行仍可能重复写入。若业务要求每批结果只生效一次,需要数据库唯一键、事务或应用自己的去重记录;Ray 的 submission ID 不会替这些外部操作去重。
6. 删除中断后继续什么,计算又由谁恢复
finalizer 是对象上表示删除前仍有待处理事项的标记。用户删除仍在运行的 RayJob 时,Kubernetes 先设置删除时间戳;只要 finalizer 尚未移除,对象就暂时保留,供 Controller 执行相应处理。它本身不包含清理代码,实际动作仍由 Controller 实现。
删除过程中还会用到 owner reference。它与 finalizer 各自负责的动作,见图 4。
RayJob Controller 看到删除时间戳后,对尚未终止的 Ray 作业尝试调用 StopJob。随后执行以下移除 finalizer 的代码:
1 | |
若进程在更新前退出,删除时间戳和 finalizer 仍在,新进程可再次进入删除分支。但此版本只记录 StopJob 的错误,随后仍会尝试移除 finalizer,并不等待 Ray 确认程序停止。更早的 Dashboard client 构造失败会直接返回;移除 finalizer 的更新失败也会重试。因此这项 finalizer 提供停止请求的处理机会,不能保证优雅停止已经完成。
owner reference 则记录 Kubernetes 对象的拥有关系。专用 RayCluster 和 submitter Job 都引用 RayJob,垃圾回收器可据此级联回收,具体顺序受删除传播策略影响。它不会替 Controller 调用 Ray Jobs API。
这里的拥有者是 RayJob 等资源对象,Operator Pod 并不拥有这次作业的整棵资源树。因此,单独重启 Operator 不会通过这条拥有关系触发集群删除。用户删除 RayJob,或 Controller 按重试、清理策略发出删除请求,才是前面讨论的资源回收入口。
正常完成后的 TTL 清理走另一条路径。TTL 是 Time To Live,在这里指完成后保留资源多久。本例根据已保存的 status.endTime 加 60 秒计算期限,重启后可以重新计算剩余等待时间;在未启用删除 RayJob CR 的环境变量、也未另配删除策略时,只删除专用 RayCluster,保留 RayJob 和 submitter Job。
到这里恢复的是作业管理流程。如果故障还带走了 head、driver 或 worker,丢失的进程内存和计算中间结果就需要另找恢复来源,RayJob 状态没有保存这些内容。Ray 自己的 _recover_running_jobs 会重新读取作业信息并恢复对非终态作业的监控;这个动作本身也没有重建用户的内存状态。
GCS 中的元数据也不能一概称为永久保存。此版本的 InitKVManager 根据配置选择内存或持久化后端,故障后的保留条件随之不同。即使作业记录还在,它保存的也是命令、状态等信息,不能替代程序将训练进度写入检查点,也不能恢复已经丢失的普通进程内存。
排查重启后的 demo-job,应先核对当前执行轮次的集群名和 submission ID,再看 Kubernetes 对象是否存在、Ray 是否仍保存对应作业记录。若已经进入新的执行轮次,还要核对前一次执行留下的外部结果。至于 Ray 元数据、task 和 actor 在各类节点故障后能恢复到什么程度,下一篇再按故障层次展开。
参考资料
- KubeRay 源码基线(6bf05eb17a3e)
- Ray 源码基线(3f785c0711b9)
- KubeRay:ray-operator/controllers/ray/rayjob_controller.go
- KubeRay:ray-operator/controllers/ray/utils/validation.go
- KubeRay:ray-operator/controllers/ray/common/association.go
- Kubernetes 对象名称规则
- KubeRay:ray-operator/controllers/ray/common/job.go
- Ray:python/ray/dashboard/modules/job/job_manager.py
- Ray:python/ray/dashboard/modules/job/common.py
- Kubernetes Finalizer
- Kubernetes 垃圾回收
- Ray:src/ray/gcs/gcs_server.cc