第 4 篇|KubeRay 恢复:Operator 重启后如何继续工作
本文最后更新于:2 天前
前言
demo-job 从创建到结束,要经过准备集群、启动提交程序、观察 Ray 作业和回收资源。Operator 如果在其中一步退出,新进程怎样判断已经完成了哪些动作,又怎样避免从头创建一套重复资源?
本文沿同一个案例走完这条流程:RayJob 新建专用 RayCluster,使用 K8sJobMode 创建 Kubernetes Job,由这个提交用 Job(submitter Job)向 Ray 提交镜像中的程序。配置 shutdownAfterJobFinishes: true、ttlSecondsAfterFinished: 60,完成后保留集群 60 秒;另设 spec.backoffLimit: 1,允许可重试的失败再执行一次。submitter 使用默认命令,不启用额外删除策略或删除 RayJob CR 的环境变量。
下面先假设只有 Operator 退出,Kubernetes 控制面、Ray 集群及已经创建的资源仍在。每个阶段都说明正常要完成什么、中断会留下什么,以及下一轮如何继续。API 暂时不可用、资源被另外删除或 Ray 进程也退出时,恢复所需的条件会变化,相关分支放在对应阶段说明。
先看整条执行流程
一次“执行轮次(attempt)”会经过多轮 Reconcile:有时这一轮只发现集群尚未就绪,等待下一轮;有时完成了资源创建,再写回新的管理阶段。Operator 重启可以发生在同一次执行轮次中。主案例中,流程判定执行失败、且重试条件满足时,会清理旧资源并开启新的重试轮次。
RayJob 用 jobDeploymentStatus 保存 KubeRay 的管理阶段,用 jobStatus 保存从 Ray 观察到的程序状态。前者进入 Running,表示提交用 Job 已经创建,后者此时还可能没有进入 Ray 的 RUNNING。
| 阶段 | RayJob 中的管理状态 | 这一阶段要完成什么 | 中断后继续所依据的信息 |
|---|---|---|---|
| 初始化身份 | New,API 中为空字符串 |
选定集群名和 submission ID,写入 Initializing |
RayJob 中保存的身份和阶段 |
| 准备集群 | Initializing |
找到或创建 RayCluster,取得可提交的访问地址 | 已保存的集群名、RayCluster 及 Service |
| 准备提交 | Initializing → Running |
找到或创建 submitter Job,写回管理状态 | 同名 Kubernetes Job 及其状态 |
| 运行与观察 | Running |
查询 submitter 和 Ray 作业,更新观测结果 | Kubernetes Job、submission ID、Ray 作业记录 |
| 完成与清理 | Complete 或最终 Failed |
按结束时间和保留策略回收资源 | 已保存的管理结束时间、仍存在的资源 |
| 失败后再执行 | Retrying → New |
删除旧资源,再清空本轮身份 | 重试状态、失败次数、旧资源是否已经删除 |
准备集群和准备提交都属于 Initializing,它们是同一状态分支中的不同步骤。用户主动删除 RayJob 是另一个入口:Controller 每轮先检查删除时间戳,有删除请求就走删除分支,不再进入普通状态分支。下文在正常完成与重试之后单独说明它。
重启后怎样重新进入协调
哪些进展还在
旧 Operator 的缓存、工作队列、局部变量和函数执行位置随进程退出而丢失。Kubernetes 中已经保存的资源、Ray 中的作业记录以及仍存活的计算进程,则各自保留原来的进展。
| 保存位置 | 保留的信息 | 新 Operator 用它做什么 |
|---|---|---|
RayJob 的 spec |
入口命令、集群配置、重试与清理要求 | 读取用户希望执行的任务 |
RayJob 的 status |
管理阶段、集群名、submission ID、失败次数和时间 | 选择本轮要进入的分支,确定查询身份 |
| RayCluster、Job、Pod、Service | 已创建资源及各自状态 | 确认资源准备和提交程序的实际进展 |
| Ray 侧 | 作业元数据、driver 和 actor 的状态、对象数据 | 提供计算进展,各类状态有各自的恢复方式 |
| 新 Operator 内存 | 重新同步的缓存和新建队列 | 组织后续协调 |
这里的 GCS 是 Ray 管理集群元数据的服务,作业信息通过它保存;driver 运行用户程序的入口,actor 是保存自身状态的计算进程。RayJob 保存的是管理进度,没有把这些进程的全部内存复制到 Kubernetes。
从读取对象到安排下一轮
新 Operator 启动监听并同步已有对象。RayJob 的 SetupWithManager 注册了 RayJob,以及它拥有的 RayCluster、Service 和 Kubernetes Job;已有对象经初始同步形成待协调请求,后续相关对象变化也能带来新的请求。即使停机期间没有新的用户操作,已有 RayJob 仍能通过这条入口重新得到处理。
一次 Reconcile 从命名空间和名称读取 RayJob。对象已经不存在,就结束这次处理;暂时读取失败,则返回错误等待重试。读到对象后,先处理删除与配置校验,再按持久化的管理状态进入相应分支。状态说明“从哪一阶段继续检查”,实际对象说明“这一阶段已经做到了哪一步”。
状态分支完成本轮判断后,主流程会先处理重试额度,再调用 updateRayJobStatus 写回变化。下面是写回调用的连续节选:
1 | |
写回成功后,普通收尾返回延时重查要求,RayJobDefaultRequeueDuration 为 3 秒;尚未就绪等等待分支也会安排后续检查。返回普通错误时,框架按错误重试规则重新入队,不能把它一概理解为固定 3 秒。完成清理、对象已删除或保持暂停等分支则可以结束本轮而不再要求定时检查。
重启后,监听和重查让对象再次进入协调;保存的身份、阶段及实际资源,则让新进程知道下一步该做什么。
图 2 把管理阶段和重试路径连起来。接下来依次看每个阶段如何利用已保存的信息继续工作。
初始化:先保存这一轮的身份
New 分支先为 RayJob 添加 finalizer。这个标记告诉 Kubernetes:删除对象前,还有 Controller 要完成的清理工作;后面的主动删除分支会用到它。接着选定本轮资源身份:status.rayClusterName 保存专用集群名,例如 demo-job-abcde;status.jobId 保存 Ray 作业的 submission ID。
用户可选填 spec.jobId,Controller 把实际采用的值存入 status.jobId,提交命令再作为 --submission-id 传给 Ray。它们在这条流程中传递的是同一个提交标识,与 driver 的内部 JobID 各有用途。初始化函数 initRayJobStatusIfNeed 只在状态中的 submission ID 为空时选择它:
1 | |
同一函数也选定 RayCluster 名称、设置本轮开始时间,并将管理状态改为 Initializing。这些先发生在内存中,随后通过前面的统一状态写回保存。New 分支到这里就结束,本轮不会继续进入创建集群的分支。
若 Operator 在这次写回之前退出,新的身份只保存在旧进程内存中,下游集群还没有沿这条路径创建。新进程仍读到 New,可以重新初始化;已经添加的 finalizer 会保留,检查到它存在后继续执行。
若状态已经保存、只是响应丢失或进程随后退出,新进程会读到 Initializing 和原来的两个标识,下一轮直接使用它们准备资源。这样,创建资源之前就有了可在重启后找回的名字。
准备集群:按同一个名字继续检查
进入 Initializing 后,RayJob Controller 调用 getOrCreateRayClusterInstance,按保存的集群名查询。找到对象就复用;确认 NotFound 才根据 RayJob 配置创建;其他读取错误返回,等待后续再读。RayCluster Controller 另外负责创建和维护这个集群的 head、worker Pod 与 Service。
假设 RayCluster 已创建,Operator 在 Pod 尚未全部准备好时退出。Kubernetes 中的 RayCluster 和已有 Pod 会保留,新 Operator 同步到它们后,RayJob Controller 继续查询同名集群,RayCluster Controller 继续按该集群的期望配置维护资源。已经运行的 Pod 在此假设下可以继续运行,管理检查则在 Operator 恢复后接上。
如果 RayCluster 的创建响应丢失,也用同一个名字重新确认。缓存尚未反映创建结果时,再次创建同类型、同命名空间的同名对象会受到 Kubernetes 的名称约束,返回“已存在”后下一轮继续读取;不会仅因为响应不确定就另选一个随机集群名。
在尚未保存 Dashboard 地址时,RayJob Controller 检查 RayCluster 是否 Ready;未就绪就把当前集群状态记录到 RayJob,再安排下一轮。就绪后,从 head Service 得到 Dashboard 地址。这个地址是提交程序和 Controller 访问 Ray Jobs API 的入口。若进程在地址写回之前退出,后续可以重新查询 Service;已经保存时则沿用该值。
这一步恢复依赖 RayJob 身份记录和对应资源仍可读取。使用已有集群的 clusterSelector 模式时,查询不到集群会返回错误,由使用者提供集群,Controller 不代建。主案例采用的是新建专用集群。
准备提交:识别已经创建的 Kubernetes Job
集群入口准备好后,同一个 Initializing 分支进入提交用 Job 的准备。createK8sJobIfNeed 按 RayJob 的命名空间和名称查询 Kubernetes Job:主案例中它也叫 demo-job。找到就返回成功,确认不存在才构造并创建。主流程随后执行下面的代码:
1 | |
这里先确保 Job 存在,再把管理状态设为 Running,最后才通过统一写回保存。创建 Kubernetes Job 和更新 RayJob 是两个 API 请求,图 3 展示其中的中断位置。
可以沿三个时间点看恢复过程:
- 尚未创建 Job 就退出:RayJob 仍是
Initializing,下一轮查询不到同名 Job,继续创建。 - Job 已保存、管理状态尚未写回就退出:Kubernetes 已可以启动提交程序,RayJob 仍是
Initializing。新进程查询到已有 Job 后复用它,再写回Running。 Running已保存后退出:新进程进入运行分支,查询现有 Job 和 Ray 作业的进展。
创建响应丢失时也先查询原名。缓存暂时未同步,重复创建会受到同名对象约束,后续协调重新读取。这里的 getOrCreate 由查询和创建两个操作组成,依靠稳定身份与重复检查接续流程。
Operator 退出不会连带停止提交用 Job。它是 Kubernetes 保存的资源,由 Kubernetes 的 Job 控制器管理其 Pod;一旦已经创建,提交程序可能在 Operator 停机期间把 Ray 程序成功提交,甚至等到它执行完毕。管理状态晚于实际进展是恢复时需要重新查询的原因。
运行:重新查询提交与计算进展
提交程序自己的重试怎样接续
此时有两条并行推进的工作:submitter 向 Ray 提交并跟踪程序,Operator 查询两侧状态。只重启 Operator 不会要求 submitter 从头运行;如果提交 Pod 自己失败,Kubernetes Job 才按提交层的重试配置再次运行提交程序。
默认命令先等待 Dashboard/GCS 健康检查通过,再用当前 submission ID 查询作业。生成命令的 BuildJobSubmitCommand 把查询放在 shell 条件中:
1 | |
jobStatusCommand 是 ray job status。查询成功就继续跟踪现有作业的日志;退出码非零时,才尝试带同一个 --submission-id 提交程序,然后跟踪日志。所以在“Ray 已接受提交、响应没有到达 submitter”的情况下,重试可以通过查询找回原来的作业。
非零退出码也可能来自网络错误,两个提交者也可能同时查询到不存在。Ray 服务端因此还会检查同一个标识是否已经登记。Jobs 服务用 supervisor 管理入口进程的启动与状态;JobManager.submit_job 在启动这个管理者之前先登记作业:
1 | |
底层通过 GCS 的内部键值存储接口(internal KV)写入记录,并设置不覆盖已有键。记录仍在时,重复提交收到“已存在”的错误。客户端预查询负责找回已有作业,服务端登记负责识别重复提交;登记之后还有进程启动步骤,程序实际执行到哪里要看后续查询。
上述行为以同一个集群、同一个 submission ID、仍保留的作业记录为依据。自定义 submitter 命令是否保留预查询,取决于其实现。换了集群、标识或原记录已删除、丢失时,需要按新的记录状态处理。
Controller 怎样把当前结果写回 RayJob
新 Operator 读到 Running 后,先处理期限等条件,再取得 RayCluster、维护访问 Service,随后通过 checkSubmitterAndUpdateStatusIfNeeded 查询提交用 Job。这个函数读取 Job 的完成或失败条件,确定 submitter 是否已经结束。之后,Controller 用保存的 submission ID 调用 GetJobInfo 查询 Ray 作业。
正常收尾同时考虑两侧:Ray 程序达到终态,并且 submitter 已结束。这里的 finishedAt 来自提交用 Job 的终态条件时间;isJobTerminal 在进入下面代码前已经按 Ray 的状态计算:
1 | |
当 Ray 和 submitter 还在运行,Controller 写回当前观测,稍后再查;两侧都满足终态条件,就把管理阶段写成 Complete 或 Failed。Ray 的 STOPPED 也可以映射为管理上的 Complete,程序是否成功仍看 jobStatus: SUCCEEDED。submitter 明确失败或期限检查也可以提前结束管理流程,不必始终等到这条正常收尾路径。
假设程序在 Operator 停机期间已经结束,新进程读到的 RayJob 仍是 Running。它重新查询 Job 和 Ray 记录,得到结束结果,再写回终态。如果进程在查询成功、写回之前再次退出,下次仍从持久化的旧阶段重新查询。结果在 Kubernetes 与 Ray 中仍可读取,就能再次计算出接下来的管理动作。
查询没有得到预期结果时,下一步取决于缺少的是哪一侧的信息。submitter Job 读取出错就返回错误,不调用初始化阶段的 Job 创建函数。Ray 作业暂时查不到、submitter 尚未结束时,K8sJobMode 继续等待;submitter 已结束而 Ray 中仍查不到作业时,则判为执行失败,不会改走 HTTP 提交。
RayCluster 的读取仍经过 getOrCreateRayClusterInstance,保留了查询不到专用集群时的创建路径。因此,同处于运行阶段,提交用 Job 的查询失败与专用集群不存在,会走不同的处理分支。
对于其他持续的 Ray 状态查询错误,Controller 会把首次失败时间保存在 status.jobStatusCheckFailureStartTime,查询恢复后清除。后续协调根据这个时间检查超时,Operator 重启不会把已保存的等待起点变成零。总执行期限使用本轮的 status.startTime;Ray 与 submitter 的终态衔接也有相应期限检查。因而,恢复后的下一步可能是继续观察,也可能是发现保存的期限已经到达,转入失败或收尾流程。
执行失败:清理旧轮次后重新开始
Operator 进程重启本身不会计作一次执行失败。前面几个阶段都在接续同一轮;只有管理逻辑判定失败后,才需要决定是否重新执行程序。
checkBackoffLimitAndUpdateStatusIfNeeded 在状态分支之后更新失败次数,并检查 spec.backoffLimit。本例设置为 1,第一次可重试失败转为 Retrying,第二次失败耗尽额度,保留最终 Failed。DeadlineExceeded 和 JobStatusCheckTimeoutExceeded 这两种原因不进入 RayJob 级重试。新的失败次数与选定的管理阶段一起写回,恢复时据此判断应进入清理还是最终收尾。
进入 Retrying 后,Controller 仍保留旧集群名和 submission ID,先删除旧 RayCluster 和 submitter Job。下面是清理分支中确认资源释放的连续节选:
1 | |
两个删除函数都先查询对象:已经 NotFound 才确认删除完成;有删除时间戳就等待;仍存在且未开始删除则发出删除请求。清理只完成一半时进程退出,新的 Operator 仍读到 Retrying 和旧身份,可以确认已经删除的那一部分,再继续处理另一部分。
等两者都确认不存在,并完成适用的调度器清理后,Controller 才清空集群名、Dashboard 地址、submission ID 以及本轮观测字段,把管理阶段改回 New。若在这次状态写回之前退出,下轮仍按旧身份查询,确认资源已删除后重复完成状态重置;写回已经成功,则下一轮开始初始化新身份。
submitter 使用后台级联删除:上级 Job 消失后,其 Pod 仍由垃圾回收器异步清理。因此,上面的条件确认的是两个上级对象已删除。新的自动生成集群名和 submission ID 随新轮次生成;显式填写的 spec.jobId 仍会采用。submitter Job 继续叫 demo-job,但新对象有新的 UID。UID 是 Kubernetes 为每个资源对象分配的唯一标识,用来区分同名的前后两个对象。
到这里可以分清两种重试:spec.submitterConfig.backoffLimit 让提交程序在当前轮次内再次运行,沿用当前 Ray 作业身份;spec.backoffLimit 则清理旧环境后重新执行整个 RayJob。已有集群的 clusterSelector 模式不允许配置正数的 RayJob 重试额度。
新的执行轮次可能再次运行用户程序。如果前一次已经向外部数据库写入一批结果,后一次仍可能再次写入。需要一次业务操作的效果恰好生效一次(exactly-once)时,数据库唯一键、事务或应用去重记录也要参与;Kubernetes 的同名资源约束和 Ray 的 submission ID 分别识别各自范围内的对象与提交。
执行结束:按保存的时间回收资源
当 Complete 或最终 Failed 被写入时,updateRayJobStatus 会设置管理结束时间 status.endTime。这与 status.rayJobInfo.endTime 中从 Ray 读到的程序结束时间是两个字段。本例清理使用前者,即管理流程记录的结束时间。
后续协调读到终态,进入清理分支,调用 handleShutdownAfterJobFinishes。它从保存的结束时间和 60 秒保留配置计算清理时刻,源码中的时间计算如下:
1 | |
TTL 是 Time To Live,在这里表示完成后保留资源的时间。尚未到期,就按剩余时间重新安排协调;已经到期,就执行清理。例如,结束时间已保存,过了 20 秒 Operator 重启,后续仍按原来的结束时间计算,不会仅因重启重新等待完整的 60 秒。实际发出清理请求的时刻还取决于 Operator 何时恢复并得到处理机会。
如果进程是在终态写回之前退出,则 Kubernetes 中还没有这次管理结束时间。下一轮需要先重新确认结果、写回终态,再从新写入的管理结束时间计算 TTL。这说明“从保存的时间继续等待”以该时间已经成功持久化为前提。
本例到期只删除专用 RayCluster,保留 RayJob 和 submitter Job。删除请求尚未发出就退出,下次继续发出;请求已经生效、响应丢失或进程退出,下次查询到 NotFound 就识别为已删除;对象带删除时间戳时,后续回收由 Kubernetes 继续推进。该清理路径在删除请求处理成功后可以结束协调,不要求 Operator 一直等待所有 Pod 退出。
如果启用了删除 RayJob CR 的环境变量或另外配置删除策略,清理对象和顺序会相应变化。使用已有集群模式时,终态流程保留用户提供的集群。本节按前言中的专用集群与默认清理方式展开。
用户删除:沿删除标记继续完成清理
用户主动删除仍在运行的 RayJob 时,Kubernetes 设置删除时间戳。初始化阶段添加的 finalizer 此时开始起作用:对象保留到该标记移除,给 Controller 留出删除前处理的机会。Controller 每轮读取后先看删除时间戳,因此重启后仍能进入这条路径,而不必先恢复普通运行阶段。
owner reference 则记录对象的拥有关系,供 Kubernetes 垃圾回收器处理下级资源。图 4 把 Controller 的停止请求和垃圾回收器的资源回收分开。
RayJob Controller 对尚未终止的 Ray 作业尝试调用 StopJob,随后执行移除 finalizer 的代码:
1 | |
若在更新成功前退出,删除时间戳和 finalizer 仍在,新进程会再次执行删除分支。Controller 记录 StopJob 的错误后仍尝试移除 finalizer,推进条件并非 Ray 已确认程序停止;更早的 Dashboard client 构造失败会直接返回,移除标记的更新失败也会重试。已经成功移除标记后,对象可以继续删除,新的 Operator 读到 RayJob 不存在时就结束处理。
专用 RayCluster 和 submitter Job 都通过 owner reference 引用 RayJob,Kubernetes 可以继续级联回收,具体顺序由删除传播策略决定。它们的拥有者是 RayJob 等资源,Operator Pod 不在这条拥有关系的顶端,所以单独重启 Operator 不会沿这条关系删除整套计算资源。停止 Ray 程序由 Controller 调用 Jobs API,回收 Kubernetes 对象由资源删除与垃圾回收继续推进。
主案例以外的恢复入口
主线之外,暂停走 Suspending 清理再进入 Suspended,恢复时回到 New;InteractiveMode 将提交交给使用者,Waiting 是等待提交标识的阶段;配置校验失败则写入 ValidationFailed。这些入口各有自己的推进条件,不能把它们都套成“重启后重新提交”。
如果退出的还包括 head、driver 或 worker,计算状态也需要恢复来源。Ray 的 _recover_running_jobs 会重新读取作业信息,为非终态作业恢复监控;它不会重建用户进程已经丢失的内存。GCS 的 InitKVManager 按配置选择内存或持久化后端,元数据能否读回取决于后端。训练进度等应用状态则需要程序从检查点恢复。第五篇继续讨论这些进程与存储层的接管条件。
小结:每个阶段都有自己的继续依据
沿 demo-job 走完一次执行,恢复依靠的信息也逐步变化:初始化先保存身份,准备阶段按身份复用资源,运行阶段重新查询提交与计算结果,重试阶段确认旧资源删除,完成阶段按保存的结束时间清理。用户主动删除则由删除时间戳和 finalizer 指明还有哪些工作要处理。
这些判断依据保存在进程之外,新 Operator 能够重新读取,管理流程就能接着推进。监听与重查提供执行机会,资源身份、对象状态、Ray 作业记录和时间字段决定本轮动作。计算状态的恢复则由 Ray、存储后端和应用共同承担。
下一篇把恢复范围从一个 Operator 进程扩展到多副本接管、Kubernetes 控制面、Ray head 与持续服务入口,继续区分每一层保存的状态和能够恢复的范围。
参考资料
- KubeRay 源码
- Ray 源码
- 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