2026-09-16 00:00:00
按第 0—10 篇依次阅读:先了解计算与部署的分工,再看 KubeRay 的集群管理和 Volcano 的批调度,最后回顾 Operator 设计。熟悉基础概念的读者可以直接从第 1 篇开始。
从分布式计算和容器部署讲起,解释 Ray、Kubernetes、KubeRay 与 Volcano 各自管理什么,以及一次作业怎样从声明走到运行。
以一次 RayJob 提交说明各类资源的用途,解释 KubeRay、Kubernetes 与 Ray 如何分工完成集群准备、程序执行和收尾。
从 worker Pod 的状态变化进入监听、缓存和工作队列,跟踪一次 Reconcile,以及事件、延时和错误怎样触发下一轮协调。
以 100 个 RayCluster 为推演场景,分析协调协程、队列键锁、缓存滞后与 expectations,区分管理并发和 Ray 计算并发。
把中断放在资源创建与状态写回之间,沿持久化身份、对象查询和重试分支,分析 Operator 重启后的恢复依据与防重边界。
从 Operator 的 Lease 选主与接管,追到 Kubernetes 控制面、Ray head、GCS 和 RayService,分别核对管理恢复、计算恢复与请求可用性。
从几个 Pod 争用资源的问题进入 Volcano,区分 Job、PodGroup、Queue,以及 Scheduler、Controller Manager 和 Admission 的协作。
沿一个三成员 PodGroup 走过准入、试分配、Gang 检查、提交与异步绑定,说明资源不足和部分绑定失败时会发生什么。
用两个 Queue 说明 Volcano 怎样计算资源份额、选择下一项工作,以及资源占满后怎样在 Queue 内抢占或跨 Queue 回收。
对比默认流程与 Volcano 接入流程,跟踪 RayCluster Controller 创建成员 Pod 的路径,再说明参数传递和生命周期协作。
总结 Operator 的资源建模与协调逻辑、并发控制、恢复和高可用设计,以及 Kubernetes 和 controller-runtime 提供的支撑机制。
2026-09-15 23:45:00
假设手里有一个 Python 程序,要处理一批数据。在自己的电脑上,装好依赖、执行 python main.py 就能开始跑。后来数据多了,希望用几台机器一起计算。这时会遇到两类问题:程序怎样利用多台机器,以及这些机器上的程序由谁部署和管理。
Ray 负责程序怎样利用多台机器,Kubernetes 管理程序的部署和运行环境。把 Ray 部署到 Kubernetes 上以后,KubeRay 又补充了哪些管理逻辑,Volcano 在什么情况下参与调度?要理解它们的关系,可以从一个 Python 程序的运行过程说起。
下面以 Kubernetes 上的 Ray 部署为例。Ray 可以在单机或多台机器上独立运行;选用 Kubernetes 和 KubeRay 后,可以先使用默认调度器,再按工作负载的需要接入 Volcano。
本文假设你运行过程序、用过命令行,不要求 Kubernetes、Ray 或 Volcano 使用经验。先解释基础概念,再走一遍从声明到程序运行的过程,最后讨论资源竞争为什么会引出批调度。集群安装和具体运维留给各项目的使用文档。
源码基线为:KubeRay 6bf05eb17a3e、Ray 3f785c0711b9、Volcano d8984501e4ad。KubeRay 的依赖按 go.mod 确定,controller-runtime 为 v0.24.1,client-go 为 v0.37.0。源码链接见文末参考资料。
先看计算本身。假设程序要处理许多文件,每个文件都可以独立完成计算,最后再汇总结果。单进程逐个处理时,后一个文件需要等前一个完成。如果希望同时利用多台机器,就得把这些工作交给不同进程执行,并把输入和结果传递过去。
可以自己写进程管理、网络通信和工作分配逻辑,也可以使用分布式计算框架。Ray 是这类框架之一:使用者通过它的编程接口表达可以分开的工作,Ray 运行时负责安排执行,以及计算对象的存储和传递。“运行时”指程序执行期间为这些操作提供支持的一组进程和机制。
Ray Core 提供 task 和 actor 两种计算抽象。task 表示一次可以在其他执行进程中异步运行的函数调用。例如,把“处理一个文件”写成 Ray 远程函数,就可以发起多个 task,再收集它们的结果。actor 则把一个带状态的对象放到执行进程中,通过方法调用使用它。例如,某个对象先加载一份模型,后续多次调用可以继续使用这份已加载的模型。
这里需要使用者或所用的上层库表达计算如何拆分。把普通 Python 脚本原样交给 Ray,并不会自动把其中的串行循环改成分布式计算。
这些工作需要一套 Ray 运行环境。Ray 集群由一个 head 节点和与它连接的 worker 节点组成:head 除了普通运行组件,还承载集群管理进程;worker 节点提供执行用户代码的进程,并参与计算调度和对象传递。head 也可以承担用户计算,具体取决于配置。运行入口脚本、发起 task 和 actor 的进程称为 driver;一次 Ray 作业包含由这个入口程序发起的计算。
Ray 可以先在一台机器上使用,也可以部署到多台机器。学习 task 和 actor 时,不需要先准备 Kubernetes。接下来引入 Kubernetes,是因为我们还希望管理这些运行环境的部署和变化。
假设已经确定要运行一个 head 和两个 worker。还需要把对应进程启动起来:选择机器,准备 Python 和依赖,配置启动参数,让 worker 找到 head。如果同时维护许多套环境,还要处理资源不足、容器退出、配置更新,以及不同程序争用机器的问题。
如果选择用 Kubernetes 管理这部分工作,就需要先把运行环境组织成容器。Kubernetes 是管理容器化应用的开源系统:使用者把可用机器组成集群,再通过统一接口描述应用需要的运行环境,由集群中的管理组件安排部署并持续观察。它可以管理 Ray,也可以管理网站、数据库等其他应用。
先解释容器。程序除了源码,还依赖解释器、软件包和系统库。容器镜像可以把程序和所需的用户态环境放在一起,供不同机器按同一份内容启动。容器是根据镜像运行起来的实例;实际启动和管理容器的软件叫容器运行时,例如 containerd。数据文件可以另外挂载或下载,不必全部放进镜像。
Kubernetes 中的 Node 是加入集群的机器,可以是物理机,也可以是虚拟机。安排到 Node 上的基本运行单位叫 Pod。Pod 包含一个或多个容器,这些容器一起部署到同一个 Node,并共享网络等资源。最简单的情况,就是一个 Pod 里运行一个应用容器。
采用 KubeRay 部署时,Ray 的 head 和 worker 节点以 Pod 的形式运行。因此,Ray 节点与 Kubernetes Node 不要求一一对应:一台机器上可以放多个 Ray Pod,只要资源和调度规则允许。一个 worker Pod 内也可以有多个执行进程,“两个 worker Pod”不等于“只能同时运行两个 task”。
两层系统作决定的对象不同。Kubernetes 为 Pod 选择机器,并负责容器运行所需的管理工作;Ray 在已经运行起来的环境中安排 task 和 actor。Kubernetes 不会因为 Python 又调用了一次远程函数,就为这个函数创建一个 Pod。
K8s 就是 Kubernetes 的缩写:保留开头的 K 和末尾的 s,用数字 8 代替中间的八个字母。两种写法指同一个项目。
kube 常出现在 Kubernetes 相关工具的名字里,例如 kubectl、kubelet。kubectl 是操作 Kubernetes 的命令行工具,可以提交声明、查询资源和查看状态;kubelet 则运行在集群节点上,后面会具体解释它怎样接手容器启动。
Kubernetes 提供的是通用部署机制。它能接收“运行这个镜像、需要这些资源”的要求,但内置逻辑并不知道一个 Ray 集群需要怎样组合 head 和 worker,也不知道提交一次 Ray 作业前应该等待哪些条件。
使用者可以自己写这些配置和管理脚本。例如,先创建 head,再创建 worker,检查集群是否准备好,然后提交程序;结束以后,再决定保留哪些资源、删除哪些资源。随着作业数量和异常情况增加,这些步骤需要持续维护。
KubeRay 将 Ray 的这部分管理逻辑接入 Kubernetes。它提供专门的资源类型,让使用者描述 Ray 集群、一次作业或一个持续服务,再由管理程序读取这些声明、创建资源并跟踪状态。它管理的对象和 Ray 调度的 task、actor 并不处于同一层。
| 希望交给 KubeRay 管理的事情 | 使用的资源类型 |
|---|---|
| 准备一套 Ray 运行环境,描述 head 和 worker 配置 | RayCluster |
| 执行一次作业,准备或选择集群,提交入口程序并处理收尾 | RayJob |
| 持续运行 Ray Serve 应用,管理集群、应用健康和访问入口 | RayService |
Ray Serve 是 Ray 上用于构建在线服务的库,可以用来提供模型推理等服务。持续服务与执行完就结束的作业有不同管理要求,所以 RayService 和 RayJob 是两种使用方式。下面以执行一次作业的 RayJob 为例。下面先看没有接入 Volcano 的基础部署,资源竞争的问题放到“资源竞争与 Volcano”一节。
把三者放在一起,一次作业的管理与计算大致经过图 1 中的交接。
图中先由 KubeRay 根据 RayJob 准备资源,Kubernetes 把对应容器启动起来。环境就绪后,KubeRay 安排提交程序把入口命令交给 Ray,Ray 再执行用户计算。后续状态检查和收尾仍由 KubeRay 负责。
“Kubernetes 启动容器”还省略了几个参与者。为了看清它们各自做什么,先暂时离开 Ray,用一个普通应用演示副本管理,再回到 KubeRay。
如果希望两个相同的应用副本持续运行,可以使用 Deployment。这是 Kubernetes 内置的一种资源类型,使用者在其中声明容器模板和副本数,相关控制器负责持续维护。这里把它命名为 intro-deployment,它与后面的 RayJob 案例无关。
在 Kubernetes API 中,Deployment 和 Pod 是资源类型,intro-deployment 是某个具体对象的名称。API 是供程序提交和查询这些对象的接口;YAML 则是表达对象内容的一种文本格式。Pod 模板描述要按什么配置创建副本,selector 则用标签条件识别这些副本。例如,给 Pod 标记 app: intro,再用相同条件选中它们。下面只摘出 Deployment 的外层结构,省略了必需的 selector 和 Pod 模板,不能直接作为部署文件使用:
1 |
|
apiVersion 和 kind 说明这份内容使用哪种资源接口。metadata 保存名称等标识;namespace 是命名空间,为对象名称提供范围,例如 default/intro-deployment。它是逻辑上的资源分组,不表示另一套机器。
spec 表达使用者的期望,这里是两个副本。许多资源还有 status,记录系统最近观察到的情况。期望两个副本时,实际可能只有一个,也可能两个都还在启动。把期望保存下来,再由管理组件持续使实际情况接近期望,就是这里所说的声明式管理。
使用者通过 kubectl 提交声明时,请求先到 API Server。它是 Kubernetes API 的服务端,接受资源的创建、查询和更新;资源数据由 etcd 保存。etcd 保存的是这类集群信息,不会自动保存 Python 程序的计算结果。API 接受 Deployment 声明之后,还需要其他组件继续工作。
Controller,中文通常叫控制器,是持续观察资源并执行管理逻辑的程序。Deployment Controller 读取副本和模板要求,管理名为 ReplicaSet(图中简称 RS)的资源;ReplicaSet 记录一组 Pod 的副本要求,由相应控制器根据缺口创建 Pod 对象。创建了 Pod 对象,只是让集群知道有这些容器需要运行。
接着由 Scheduler,也就是调度器,为尚未分配节点的 Pod 选择 Node。它需要考虑资源需求和调度约束,再通过 API 确定 Pod 与 Node 的分配关系,这一步称为绑定;如果没有合适的节点,Pod 就需要继续等待。Scheduler 选定机器以后,容器仍要由那台机器上的组件启动。
这里的资源需求来自配置。例如,容器的 resources.requests 可以声明需要多少 CPU、内存,调度器用这些请求判断节点能否放下它;resources.limits 则约束容器使用相应资源的上限。请求值与程序此刻的实际用量可能不同。
每台运行 Pod 的节点上都有 kubelet。它持续观察分配给本机的 Pod,按照其中的配置调用容器运行时,并向 API Server 回报 Pod 状态。容器运行时负责实际运行容器。这样,“哪个 Pod 归本机运行”和“本机的容器实际怎样了”就能在节点与管理端之间对应起来。
Pod 的状态还需要分开看。Pending 表示启动尚未完成,可能在等节点,也可能已经有节点、正在拉取镜像;Running 表示已进入容器运行阶段,不能单独证明应用可以接收工作。Ready 是另一项就绪条件,结合容器就绪状态等因素判断;如果配置了 readiness probe,就绪探针会检查应用是否准备好。后面 KubeRay 对集群是否就绪的判断,还会再组合多个 Pod 的情况。
API Server、调度器和控制器等承担集群管理工作的组件合称控制面。图 2 把它们与运行应用的 Node 分开画,表示职责分工,并不要求每个框都占一台独立机器。
前面的步骤靠各组件观察资源变化接续,图 3 按发生顺序把它们展开。
提交终端只负责发出请求,后续工作交给集群里的进程。若一个受管理的 Pod 被删除,ReplicaSet Controller 观察到副本数量不足,会尝试创建替代对象;新 Pod 再经过调度和启动。控制器反复比较期望与实际情况,决定下一步操作,这个过程称为协调,英文是 reconcile。
只创建一个孤立的 Pod,就没有 Deployment 和 ReplicaSet 这一层副本维护逻辑。机器故障后能否补建,还取决于故障是否被识别、相应控制器能否工作,以及集群是否有资源。即使替代 Pod 已经启动,也不能由此推断原程序的内存和计算进度已经恢复。
程序还需要相互访问。Pod 会被替换,地址也可能改变,因此常用 Service 提供相对稳定的访问入口。Service 可以通过标签选择后端 Pod,具体转发由集群网络实现。标签就是对象上的键值标记,用来识别一组资源。后面出现的 head Service 就承担类似职责,让提交程序等组件能够找到 Ray head。
回到 Ray。Deployment 能表达按照模板维护副本,Ray 集群却还有 head、worker 组、启动参数和连接关系。KubeRay 需要让 Kubernetes 保存这些额外信息,也需要有人按它们执行管理动作。
增加资源类型的机制之一叫 CRD,全称 CustomResourceDefinition,即自定义资源定义。它向 Kubernetes 声明一种资源叫什么、有哪些字段、如何校验。KubeRay 仓库里就有 RayCluster 的 CRD 文件,下面摘出名称定义,省略同级的其他字段:
1 |
|
这段内容让 API 知道资源属于 ray.io 这个组,类型名是 RayCluster,并按 namespace 区分命名范围。安装完整 CRD 后,再创建一个名叫 intro-cluster 的 RayCluster,这个具体对象就是自定义资源,常简称 CR。CRD 定义类型,CR 是这种类型的实例。intro-cluster 是这里为说明概念选取的示例名称。
CRD 让 API 能保存集群声明,具体管理动作由 KubeRay Operator 执行。它读取 RayCluster,根据 head 和 worker 配置创建、维护 Pod 及 Service,并回写集群状态。只有资源定义、没有相应控制器时,API 不会自行推导出这些步骤。
把应用的管理逻辑写进控制器,通过自定义资源驱动它工作,就是这里使用的 Operator 模式。KubeRay Operator 通常也作为 Pod 部署在集群中,它通过 API 读写资源,不需要成为 API Server 进程的一部分。
这几个概念可以这样区分:Operator 是部署起来的管理程序,Controller 是其中负责某一类资源的控制逻辑,Reconcile 则是控制器处理一个待协调对象时调用的函数。KubeRay 的一个 Operator 进程可以容纳多个 Controller,三类核心资源并不意味着需要三个独立进程。
现在可以把图 1 中“准备 Ray 运行环境”展开到 Pod 和进程。图 4 假设声明需要一个 head Pod 和两个 worker Pod,省略上层 RayJob 和提交程序,先看这套集群怎么运行起来。
KubeRay 读取 RayCluster 并创建所需资源;Scheduler 给 Pod 选择 Node;节点上的 kubelet 和容器运行时启动 Ray 容器。图中把一个 head Pod 和两个 worker Pod 放在两台机器上,实际位置由资源需求和调度规则决定。Operator 也可以运行在这些 Node 上,为突出职责,图中单独列出它。
容器里的 Ray 进程组成集群之后,还需要把入口程序交给它运行。这里的 RayJob 使用专用集群,会等待集群就绪,再创建一个 Kubernetes Job 来运行提交程序。Kubernetes Job 是管理运行至完成的 Pod 的内置资源;这里的 Pod 运行 Ray CLI,也就是 Ray 的命令行工具,通过 Ray Jobs API 提交入口命令。
因此,一条流程里会同时出现 RayJob、Kubernetes Job 和 Ray 作业。RayJob 描述这次作业的管理要求;提交用的 Kubernetes Job 管理提交程序的 Pod;Ray 接收入口命令后运行用户程序,安排它发起的 task 和 actor。资源创建成功、提交程序启动、用户计算完成,是三个不同的进展。
计算开始以后,Ray 节点之间的计算数据不经过 Operator。KubeRay 继续检查作业状态,并按声明决定怎样收尾。Pod 是否存在、容器是否就绪、程序是否成功,各有对应的观察依据。
程序及依赖怎样到达容器,也需要提前确定。这里假设它们已经包含在运行镜像中。声明一个 RayJob 本身不会把客户端电脑上的任意目录上传到集群;本地打包、上传和持久保存,需要对应的提交或制品管理流程。
前面的步骤解决了怎样创建和维护 Ray 集群。若同时提交几套作业,机器资源只能满足其中一部分,接下来还要决定谁先运行,以及怎样在作业之间分配资源。
例如,两个作业都至少需要三个成员才能有效计算,机器却只能容纳四个等大的成员。如果各自启动两个,双方都可能占着资源等剩下的成员。默认的 Pod 调度路径不会自动从业务代码推导出这种组内依赖,需要另外表达“这一组至少满足什么条件,启动才有意义”。这只适用于有相应要求的程序;有些 Ray 程序能够利用现有 worker 逐步计算,不必等待全部成员到齐。
Volcano 是建立在 Kubernetes 上的批处理与调度系统。这里的批处理主要指一批提交后需要安排资源执行的计算工作负载。它提供成组调度、队列和资源共享策略:成组调度通常称为 Gang,检查一组成员是否满足最低运行要求;队列把采用同一资源策略的工作负载组织起来,用于决定它们怎样共享资源。
接入时,KubeRay 继续负责 Ray 集群和作业生命周期,把组的调度要求交给 Volcano。Volcano 为相关 Pod 选择节点,kubelet 仍在节点上启动容器,Ray 继续在运行环境中安排 task 和 actor。安装 Volcano 不会让计算结果自动持久化,也不会代替 KubeRay 管理作业提交和收尾。
是否接入 Volcano,要看业务有没有最低成员要求、队列间共享或优先级等调度需求,以及现有调度方式能否满足。
到这里,四个项目可以按各自处理的对象区分:
| 项目 | 在上述部署中负责什么 | 采用条件 |
|---|---|---|
| Ray | 执行用户程序,安排 task 和 actor,管理计算对象 | 需要分布式计算,可以独立部署 |
| Kubernetes | 保存部署声明,安排 Pod 到节点,管理容器运行 | 选择用 Kubernetes 管理容器部署 |
| KubeRay | 根据 RayCluster、RayJob 等声明维护 Ray 资源和生命周期 | 选择用 Operator 管理 Kubernetes 中的 Ray |
| Volcano | 根据组需求和队列策略,为相应 Pod 分配资源与选择节点 | 按调度需求接入 |
2026-09-15 23:40:00
第零篇解释了 Ray、Kubernetes 和 KubeRay 的基本分工。这一篇把分工放进一次 RayJob 提交:使用者只创建一个 RayJob,集群里随后却出现 RayCluster、head Pod、worker Pod,以及一个运行提交程序的 Kubernetes Job。这些对象分别由谁创建,程序又从哪一步开始交给 Ray?
本例使用名为 demo-job 的 RayJob 运行 Python 程序。KubeRay 为它准备一套专用集群,其中有一个 head Pod 和两个 worker Pod;集群就绪后再提交程序,成功完成后保留集群 60 秒。程序及依赖已经放进运行镜像,暂不启用自动扩缩容。
下文把提交用的 Kubernetes Job 称为 submitter Job,把它创建的 Pod 称为 submitter Pod。第一至第五篇沿用这个不启用 Volcano 的基础案例;第六至第八篇单独介绍 Volcano,第九篇再回到 demo-job 说明接入方式。本文先走完资源创建、程序提交和收尾,队列、恢复与多副本接管留到后文。
沿用第 0 篇中的分工:Kubernetes 将 head 和 worker Pod 部署到节点上;Ray 在容器内安排 task 和 actor;KubeRay Operator 则根据 RayJob 等声明准备资源、跟踪状态。下面把这层分工落到两个具体接口上:Kubernetes API 管理资源对象,Ray Jobs API 接收入口命令并提供作业状态查询。
顺着图中的连接看,Operator 一侧通过 Kubernetes API 管理资源,另一侧通过 Ray 管理 API 查询作业或管理 Serve 应用。Ray 节点之间传输对象、执行 task 和 actor,不需要经过 Operator。
Operator 的主体是 Go 程序。Ray 运行时和提交程序分别运行在其他容器里,后面会展开它们的位置。
使用 KubeRay 时,最常见的入口是 RayCluster、RayJob 和 RayService。
CRD 定义自定义资源的结构,CR 是它的具体实例。安装 RayJob CRD 以后,就可以创建 demo-job 这样的 RayJob 对象。Kubernetes 保存对象,Controller 读取声明并执行管理逻辑。
选择哪种资源,取决于希望 KubeRay 帮忙管理到哪一步。
| 资源 | 使用者声明什么 | KubeRay 主要管理什么 |
|---|---|---|
| RayCluster | 一套 Ray 集群的配置,包括 head 和各组 worker | 集群所需的 Pod、Service,以及集群状态 |
| RayJob | 一次作业的入口命令、运行环境、集群来源和收尾要求 | 准备集群、提交作业、跟踪执行、按策略清理 |
| RayService | Ray 集群配置与 Ray Serve 应用配置 | 承载应用的集群、Serve 部署、健康检查和访问入口 |
直接创建 RayCluster,可以得到一套供程序连接的 Ray 集群。RayJob 还要负责提交一次作业,跟踪结束状态并处理清理。RayService 面向持续提供服务的 Serve 应用,关注应用健康和访问入口。这两种上层资源都能复用 RayCluster 的集群管理逻辑。
图中的三列是独立的使用方式,分别对应不同的 RayCluster;连线表示资源的创建和管理关系,辅助资源有所省略。下文沿中间一列的新建专用集群路径阅读。
以 RayJob 新建专用集群为例,它的 spec.rayClusterSpec 中包含集群配置。RayJob Controller 会根据这份配置构造一个独立的 RayCluster 对象。下面是 constructRayClusterForRayJob 中的连续节选:
1 |
|
DeepCopy() 把 RayJob 内嵌的集群配置复制给新对象。Labels 和 Annotations 都是对象上的键值元数据:label 常用于筛选资源,annotation 常用于附加说明或组件约定的配置,具体含义由读取它的程序解释。随后设置的 owner reference 记录 RayCluster 属于这个 RayJob,供 Kubernetes 垃圾回收等机制使用;这里的垃圾回收指依据资源拥有关系清理下级对象。新对象提交到 API 后,具体的集群创建工作就由 RayCluster Controller 接手。
RayJob 也支持通过 clusterSelector 使用已有集群,此时不创建新的专用集群,收尾时也要遵循已有集群的生命周期约束。
KubeRay 核心 Operator 的一个实例运行一个 Go 主进程。controller-runtime 是用于编写 Kubernetes 控制器的 Go 库,其中的 Manager 对象负责组织 Controller 的运行,共享客户端、缓存等基础设施。RayCluster、RayJob 和 RayService 各有自己的 Controller,但它们注册在同一个 Manager 上。
main.go 中的三处注册如下。这里拼接了三段原始调用,省略中间的选项构造和其他可选组件注册。
1 |
|
这三个调用传入的是同一个 mgr,所以三个 Controller 都运行在同一个 Operator 进程里。增加 Operator 副本时,增加的是整套进程实例。
图 3 按 Pod 分组,列出容器中的主要进程或逻辑模块。实际 Ray 进程数会随配置和工作负载变化。Service 提供网络访问入口,属于资源对象,不在进程列表中。
head Pod 的 Ray 容器里包括 GCS、Dashboard 相关进程和 raylet。GCS 全称 Global Control Service,管理集群级元数据;Dashboard 除了页面,还提供 Jobs、Serve 等管理 API。raylet 参与节点上的任务调度和对象管理。worker Pod 里也有 raylet,用户的 task 和 actor 则在相应执行进程中运行。head 节点是否承担用户计算,取决于它的资源配置。
这部分结构可以对照 Ray 的 start_head_processes 和随后的 start_ray_processes 阅读。KubeRay 的 common/pod.go 负责构造 Pod、准备 ray start 启动命令,Ray 再负责启动自己的运行时进程。
如果启用 Ray 自带的自动扩缩容功能,即 in-tree autoscaling,head Pod 会增加运行扩缩容程序的 autoscaler 容器。这里的 in-tree 指扩缩容程序随 Ray 集群部署。它观察 Ray 的资源需求并调整 RayCluster 的扩缩容目标,RayCluster Controller 据此管理 worker Pod,Kubernetes 负责调度新增 Pod。Ray 侧的 KubeRay node provider 负责把扩缩容决定转换为 RayCluster 配置更新,可以从这里继续查看实现。
本例中的 submitter Pod 运行 Ray CLI,负责向 head 上的 Jobs API 提交命令。用户入口程序随后由 Ray Jobs 机制安排执行,它的运行位置与 submitter 是两回事。
Controller 会反复执行一段称为 Reconcile 的协调逻辑。每次处理某个资源时,它读取声明和当前观察到的状态,判断还缺什么,再执行相应动作。
例如,配置要求两个 worker,RayCluster Controller 当前只看到一个,就可以创建缺少的 Pod。API 接受创建请求以后,还要由 Kubernetes 调度并启动容器,Controller 后续再检查 Pod 状态。这些步骤会分多次完成。
在资源对象中,spec 表达期望,status 记录 Controller 观察和处理后的状态。Pod、Service 等实际对象也是后续判断的依据。Controller 每轮都要重新检查,不能仅凭“上次已经发出创建请求”就认定环境准备完成。
RayJob Controller 不只关注 RayJob 本身,也关注与它有关的下级资源。在 SetupWithManager 的开头,可以看到这些声明:
1 |
|
For 指定主要处理的资源类型。Owns 声明对相关从属资源的关注,这些资源发生事件时,可以触发其所有者的协调。事件如何映射回某个 RayJob、如何进入队列,下一篇再顺着这段代码往下读。
RayJob Controller 与 RayCluster Controller 通过 Kubernetes 资源及其状态交接工作:前者创建集群声明,后者创建 Pod、观察就绪情况并更新集群状态,前者再读取这个结果。整个过程跨越多轮 Reconcile。
把前面的分工放进一个具体例子:运行镜像里已经准备好 Python 程序 /app/main.py 及其依赖,我们希望用一套专用 Ray 集群执行它,结束后释放集群资源。
作业的名字叫 demo-job,放在 Kubernetes 的 default 命名空间中。命名空间用于区分资源的命名范围;后文写成 default/demo-job 时,斜杠前面是命名空间,后面才是对象名。
集群由一个 head Pod 和两个 worker Pod 组成。worker group(worker 组)把采用相同 Pod 模板和资源配置的 worker 放在一起配置。本例只有一组,自取名为 compute,意为“计算”。组的副本数设为 2,就表示本例需要两个 worker Pod;多 host 配置下的数量关系留到第九篇说明。
提交模式选择 K8sJobMode:由一个 Kubernetes Job 运行提交程序,把入口命令交给 Ray。下面只摘出与提交、清理有关的 YAML,必填的集群和 Pod 模板已省略,片段不能直接提交。
1 |
|
RayJob 描述一次作业的管理要求,RayCluster 描述这次作业使用的集群。二者是两种 Kubernetes 资源类型,各有自己的对象名。本例的 RayJob 叫 demo-job,由它创建的 RayCluster 示意为 demo-job-abcde。末尾的 abcde 是 Controller 自动生成的随机后缀。
提交程序把入口命令交给 Ray Jobs API 时,还会携带一个 submission ID,即“作业提交标识”。Ray 用它关联这次提交的状态和日志。它标识的是 Ray 中的一次作业提交,和 Kubernetes 中的集群对象名用途不同。本例没有指定 spec.jobId,Controller 会生成一个值,示意为 demo-job-fghij。
下面把这些名称放在一起对照。abcde 和 fghij 是两次独立生成的后缀示意,实际集群名和提交标识保存在 RayJob 的 status 中。
| 对象或标识及其用途 | 本例中的值 | 名称由谁确定 |
|---|---|---|
| RayJob:描述作业如何提交与收尾 | default/demo-job |
使用者在 YAML 中指定 |
| RayCluster:描述专用 Ray 集群 | default/demo-job-abcde |
Controller 生成,保存到 status.rayClusterName
|
| worker group:把同配置的 worker 放在一组 | compute |
使用者在 rayClusterSpec 中命名,本例配置两个 worker Pod |
| submitter Job:让 Kubernetes 运行提交程序 | default/demo-job |
沿用 RayJob 的名称;资源类型不同,所以可以同名 |
| submission ID:让 Ray 识别这次作业提交 | demo-job-fghij |
Controller 生成,保存到 status.jobId
|
集群名和 submission ID 各自生成随机后缀,Controller 分别保存它们的对应关系,不能靠一个推导另一个。status.jobId 会作为 Ray CLI 的 --submission-id 参数传给 Ray Jobs API。
Ray 内部还会给 driver 分配 JobID。driver 是执行入口程序、发起 Ray 计算的进程;这个内部 JobID 与前面的提交标识不是同一套编号。submitter Job 创建的 Pod 也有自己的名称,因为 Job 是管理对象,Pod 才是运行提交程序的单位。第四篇讨论整次执行重试时,还会看到集群名和 submission ID 重新生成。
本例通过镜像提供程序和依赖。也可以挂载文件,或配置 Ray 的 runtime_env 来分发代码。runtime_env 描述作业的运行环境,可以指定代码、Python 依赖等内容。
图中省略事件传递和重复检查,只保留准备集群的关键步骤。“Kubernetes”一列合并表示调度器、kubelet 和容器运行时等执行组件。
RayJob Controller 首先初始化执行所需的身份信息,例如 submission ID 和 RayCluster 名称,并把状态推进到 Initializing。接着,getOrCreateRayClusterInstance 根据名称查找集群;如果专用集群尚不存在,就创建它。
这一步创建的是 RayCluster CR。接到这份集群声明以后,RayCluster Controller 才会去维护 head Service、head Pod 和 worker Pod。相关入口是 reconcileHeadService 和 reconcilePods。
Pod 被创建以后,容器启动交给 Kubernetes。RayCluster Controller 继续观察资源状态,把集群状况写回 RayCluster 的 status。RayJob Controller 后续读到这个状态,才知道集群准备到了哪一步。
等待条件在 RayJob 的 Initializing 分支中。下面保留条件和状态赋值,省略日志、注释及 URL 查询代码:
1 |
|
集群尚未就绪时,break 退出当前状态分支,后续代码更新状态并安排再次检查。等到观察到集群就绪,Controller 才会获取 Dashboard 地址并推进提交。这里的 Go 常量 rayv1.Ready 对应 RayCluster 的 status.state: ready,与 Pod 的 Ready 条件、RayJob 的 Running 阶段分别属于不同资源和字段。到这一步,准备好的是运行环境,入口程序还要通过 Jobs API 提交。
下面接着看提交、执行和清理。图中的 TTL 指完成后保留资源的等待期限,本例设为 60 秒。
集群准备好后,沿图 5 继续看一次成功的提交和收尾。图中的 Kubernetes 一列合并了 Job 控制器、调度和节点执行环节。
RayJob Controller 在 K8sJobMode 分支中创建提交用的 Kubernetes Job。下面是对应代码的连续节选:
1 |
|
Kubernetes 自己的 Job 控制器根据这个 Job 创建 submitter Pod。submitter 运行 Ray CLI,通过 head Service 访问 Dashboard 的 Jobs API。
BuildJobSubmitCommand 负责拼接提交命令。把健康检查、已有 job 检查和日志跟踪等外围逻辑拿掉,其核心命令形态如下。这是带占位符的示意,不是完整生成结果:
1 |
|
这条提交命令带有 --no-wait,完整脚本随后还会跟踪作业日志。因此,提交命令返回之后,submitter 容器仍可能继续运行。
Ray 接受请求后,通过 Jobs 机制启动和管理入口程序。入口程序使用 Ray API 发起 task 或 actor,由 Ray 运行时安排计算。
与此同时,RayJob Controller 会调用 Dashboard client 的 GetJobInfo,查询 Ray 侧的 job 状态,并结合 submitter Job 等信息推进 RayJob 状态,最终通过 status 更新写回 Kubernetes API。
执行状态也沿着这套分工分别记录。三种状态描述的是不同层次的进展:
| 对象 | 由谁管理 | 状态主要回答什么 |
|---|---|---|
| RayJob CR | KubeRay 的 RayJob Controller | 集群准备、作业提交和收尾流程走到哪一步 |
| submitter Kubernetes Job | Kubernetes 的 Job 控制器 | 提交程序对应的 Pod 执行得怎样 |
| Ray 运行时中的 job | Ray Jobs 机制 | 用户入口程序是否等待、运行、成功、失败或停止 |
RayJob 的 status 中也区分了 jobDeploymentStatus 和 jobStatus。前者描述 KubeRay 管理流程,后者反映 Ray 作业状态。上面的源码在创建提交用 Job 后就把 jobDeploymentStatus 设为 Running,因此这个值不能单独证明用户程序已经进入 Ray 的 RUNNING 状态。
判断程序是否成功,要继续看 Ray 作业结果和 RayJob 后续状态。kubectl apply 成功说明声明已经提交,submitter Pod 启动说明提交程序开始运行,二者都还没走到用户程序完成这一步。
本例设置了 shutdownAfterJobFinishes: true 和 ttlSecondsAfterFinished: 60,使用这一收尾路径,且没有启用删除 RayJob CR 的环境变量。对于正常完成的作业,Controller 按结束时间计算等待期限,到期后删除专用 RayCluster,再由 Kubernetes 按资源拥有关系回收下级资源。
handleShutdownAfterJobFinishes 在这条路径中保留 RayJob CR 和提交用的 Kubernetes Job。这样还能查询 RayJob 的状态,并在相关 Pod、日志仍可访问时查看提交日志。
结果文件、日志和程序版本仍需另外归档。保留下来的 RayJob 对象主要记录这次执行的声明与状态。
一次 RayJob 提交把几层管理工作接了起来:RayJob Controller 安排集群准备、程序提交与收尾;RayCluster Controller 维护运行环境;Kubernetes 调度 Pod、启动容器;Ray 接手入口程序并安排计算。各层通过资源声明、状态和管理 API 交接,运行中的 task 和 actor 不需要经过 Operator。
RayCluster 被单独设计成一种资源,是因为集群的创建与维护可以服务于不同的使用方式。上层入口改变时,底层维护 head、worker 的逻辑仍可复用。
| 使用需求 | 选择的入口 | 集群之上的工作由谁负责 |
|---|---|---|
| 准备一套可长期使用的 Ray 集群 | RayCluster | 使用者自行安排作业提交 |
| 提交一次作业,并按策略收尾 | RayJob | RayJob Controller 管理提交、状态与清理 |
| 持续提供 Ray Serve 服务 | RayService | RayService Controller 通过 Serve API 管理应用、检查健康并维护入口 |
RayService 需要替换集群时,还会区分 active(当前服务集群)和 pending(准备接替的集群)。它先准备接替者,再调整服务入口,对应实现位于 rayservice_controller.go,第五篇会展开这种设计的条件与边界。
下一篇继续看资源变化怎样触发协调:一个 worker Pod 的事件,如何映射到对应的 RayCluster,再进入 Controller 的工作队列。
本文 Go 节选仅调整缩进;拼接和省略范围已在上下文标明。
2026-09-15 23:35:00
第一篇里,RayJob default/demo-job 创建了专用 RayCluster,示例名称为 default/demo-job-abcde,包含一个 head Pod 和 compute 组的两个 worker Pod。RayJob Controller 要等集群就绪,才会创建 submitter Job。假设最后一个 worker Pod 已经进入 Running,Ready 条件也刚变为 True,这个变化怎样传到 Operator,又怎样让作业继续提交?
沿着这次 Ready 变化,可以把协调过程分成三段:事件先把 RayCluster 标识放进队列;一轮 Reconcile 维护资源、计算状态,再把结果交给 RayJob Controller;本轮返回以后,新的事件、延时请求或错误重试继续安排下一轮处理。
第 0 篇中负责回报 Pod 状态的 kubelet,会把本节点观察到的结果写入 API Server。Operator 通过 watch 接收所关注资源的变化通知;这是 Kubernetes API 提供的监听机制。
例如,某个 worker 的 Ready 条件从 False 变成 True,API Server 中相应 Pod 的记录也随之更新。监听这个 Pod 的客户端可以收到更新通知,得知这个对象发生了变化。通知提供了重新检查资源的机会,至于整个集群是否已经就绪,还要由后面的协调逻辑判断。
管理这些对象同步和通知的客户端组件叫 informer。它在 Operator 进程内维护资源的本地缓存:例如某个 Pod 的名称、所属集群和已观察到的状态,都可以从这份对象记录中读取。收到对象更新后,informer 更新本地记录,并通知注册的事件处理器,源码中通常称为 handler。
handler 负责确定这次变化应该让谁重新检查。在本例中,发生变化的是 worker Pod,需要重新检查的是它所属的 RayCluster。因此,处理器要从 Pod 找到所属集群,再把这个集群的标识放进待协调队列。这个处理器负责安排检查,集群是否需要补 Pod、能否标记就绪,留到 Reconcile 中判断。
Controller 的运行框架在后台从队列取出标识,调用 Reconcile;需要的资源对象则从客户端读取。负责取队列、调用函数的 goroutine 在本文中称为协调协程。goroutine 是 Go 运行时调度的执行单元,可以在同一进程内共享内存,与执行用户计算的 Ray worker 各有职责。
KubeRay 在 main.go 中创建 Manager,配置缓存、监听的 namespace 和客户端。Manager 是 Operator 进程里的框架对象,各个 Controller 通过它注册监听,共享缓存和客户端等基础设施。这些组件在同一管理程序中协作,不需要分别部署成独立服务。
这里的“监听”有范围。WatchNamespace 决定关注哪些 namespace;KubeRay 还在 internal/managercache/cache.go 中为 Pod、Kubernetes Job 等对象设置标签筛选。普通业务 Pod 的每次更新没有必要都送到 KubeRay。监听范围之外还要看权限:RBAC 是 Kubernetes 的基于角色的访问控制,用来决定这个 Operator 身份能对哪些资源执行哪些操作。配置了监听范围,不等于已经获得读取权限。
Controller 开始处理对象之前,需要先完成缓存的初始同步。controller-runtime 的 Kind.Start 注册事件处理器后,会等待缓存同步,并等待这个处理器收到初始事件。Controller 启动协调协程之前,还要等待这些事件源同步完成。因此,启动时 watch 权限不足、资源类型不可用或缓存同步失败,都可能让 Controller 无法进入正常处理阶段。
RayCluster Controller 注册监听的代码位于 SetupWithManager。这段代码是在声明监听规则:哪些资源发生变化时,需要让这个 Controller 再检查一次 RayCluster。其中 predicate 是事件筛选条件,generation 是用于识别期望配置变更的计数。下面保留构建器调用,省略后面的调度器配置和 Controller 选项:
1 |
|
For 指定这个 Controller 主要协调 RayCluster。RayCluster 自身发生符合条件的事件时,处理器把它的 namespace 和 name 放进队列。Owns 关注下级资源:Pod 变化时,处理器沿它的 owner reference 找到所属 RayCluster,再把这个集群的标识放进队列。
对照代码中的括号,WithPredicates(...) 是传给 For(RayCluster) 的参数,只筛选 RayCluster 自身这条监听路径。三个条件用 predicate.Or 连接:对更新事件,只要 generation、label 或 annotation 有一项变化,就可以通过这组筛选。
例如,使用者修改 worker 数量要求,RayCluster 的期望配置发生变化,generation 随之变化,可以触发重新检查;Controller 只把观察结果写回 status 时,通常不会改变 generation,也就不会仅凭这一条件触发自身协调。这样可以避免每次写回状态都立即通过同一条监听路径再次入队。
代码中的 Owns(Pod) 没有传入这组筛选条件,因此 worker 的 Ready 条件变化仍然可以触发集群检查。RayJob 对下级 RayCluster 的监听也独立配置。同一次 RayCluster 状态更新,可以被 RayJob Controller 接收,而被 RayCluster Controller 自身的这组条件过滤掉。
owner reference 是创建资源时写入的拥有关系,Owns(Pod) 利用它注册下级资源的监听。比如最后一个 worker 就绪时,处理器沿这份关系找到 demo-job-abcde,把这个 RayCluster 放入队列,后续检查便能同时考虑 head 和其他 worker。
这段注册还包括 Secret 和 PersistentVolumeClaim(PVC)。前者用于存放密码等敏感配置,后者用于声明持久存储需求。不同集群配置会用到不同辅助资源,列在 Owns 中不表示每套 RayCluster 都会创建它们。
这个行为能在 controller-runtime 的 Builder.doWatch 中看到:For 使用 EnqueueRequestForObject,Owns 使用 EnqueueRequestForOwner;默认情况下,后者只匹配标记为 controller owner 的引用,也就是 owner reference 中 controller: true 的那一项。对象可以记录多个 owner,但至多有一个这样标记的管理者。处理器检查直接拥有关系,不会递归寻找任意层级的祖先。
owner handler 匹配资源的 API 组 Group 和类型 Kind,例如 ray.io 与 RayCluster,再构造下面的请求。以下是 getOwnerReconcileRequest 的连续节选:
1 |
|
后续代码根据资源作用域补上 namespace。对 RayCluster 来说,请求标识形如 default/demo-job-abcde,只告诉协调函数该检查哪套集群。至于要不要创建 worker,要等读到对象后再判断。
前面的 For 和 Owns 接收不同来源的事件,最后都可以得到同一个 RayCluster 标识,两条入口由此汇合到同一套协调逻辑。RayJob 对下级 RayCluster 也有自己的监听,形成图中的两层关系。图里的每一层都先把所属对象入队,再由对应 Controller 处理;实际的状态交接要等协调函数执行以后才会发生。
假设 compute 组的两个 worker 接连更新,处理器从它们的拥有关系都找到了 default/demo-job-abcde。如果每条通知都对应一次独立协调,Controller 就可能连续检查同一套集群,重复读取相同的资源。队列因此按对象标识合并尚待处理的请求:在取得处理机会之前,同一个 RayCluster 的多次入队可以汇成一个待处理项。
这里可以对照缓存和队列各自保留的内容。缓存保存资源对象,例如各个 worker 已观察到的状态;队列中的 default/demo-job-abcde 只表示“这套集群还需要检查一次”。它不记录每条 Pod 通知,也不携带“创建一个 worker”这样的操作指令。只要下一轮重新读取集群和所属 Pod,就可以一起判断这些变化带来的结果。
取得处理机会以后,情况又有区别。假设这一轮已经读过 Pod,另一个 worker 才更新,新的请求就有必要留下。队列允许为正在处理的对象登记下一轮检查;在本轮结束前,同一队列不会把这个对象再次交给另一个协调协程。期间又到来的相同请求仍可合并。等本轮完成、下一轮开始后,再有变化,还可以继续安排检查。
因此,去重针对的是尚待处理的请求,没有把这个对象永久标记成“已经处理过”。队列也无法判断本轮是否恰好已经读到了新变化:有时下一轮会发现资源已经符合期望,无须再写入。这要求协调逻辑能够反复比较期望和当前状态,写入操作也要考虑重复调用。幂等要求重复执行同一操作不会额外产生一次效果;控制循环本身不会替被调用的外部接口实现这个要求。
协调协程取得这份请求后,先用其中的标识读取 RayCluster,再检查相关资源。触发请求的 worker 可能在进入 Running 后又将 Ready 条件更新为 True,也可能已经删除;本轮依据当前能读到的对象作判断,无须重放中间每一条通知。
“写 API 成功”和“缓存里已看到新对象”之间存在时间差。Manager 提供的常规类型客户端通常从缓存读,写操作直接发给 API Server;绕过缓存的读取另有路径。这个差异在第三篇分析并发创建时还会用到。
RayCluster 的 Reconcile 开头就进行了读取。下面是连续节选:
1 |
|
这段代码用请求里的 namespace/name 读取 RayCluster,把结果放进 instance,读取成功才进入实际维护函数。后面的 NotFound 分支处理对象已删除、旧请求仍留在队列里的情况,此时可以结束本轮;连接或权限等读取错误则要按失败处理。
进入 rayClusterReconcile 后,代码会检查资源是否交给外部 Controller 管理,验证配置,处理删除流程,再按当前状态维护集群资源。不同配置会走不同分支。本例没有删除请求,也不启用自动扩缩容,重点是 head Service、head Pod、worker Pod 和集群状态。
spec 保存使用者的期望,status 保存 Controller 已观察到的集群情况;实际 Pod、Service 等对象也参与判断。只读某一个 status 字段,不能代替完整的资源检查。
假设配置要求两个 worker,本轮只观察到一个,Controller 会结合已有 Pod 和尚未被缓存观察到的创建记录,判断是否需要补建。创建请求发出后,Pod 可能还要等待调度和镜像拉取,本轮只完成当前能做的管理动作;如果数量已满足,就继续检查这些 Pod 的状态。
资源维护之后,还要把观察结果整理成集群状态。在 raycluster_controller.go 中,rayClusterReconcile 执行完各项维护操作,接着调用负责计算状态的 calculateStatus:
1 |
|
这里的 instance 是当前 RayCluster,reconcileErr 记录前面资源维护是否出错。calculateStatus 复制这个对象,列出所属 Pod,再把计算结果放进返回的 newInstance;调用方在计算成功后将新状态写回 API。计算状态和写回状态是这一步里的两项操作。
本例关注集群第一次进入 Ready 的条件。下面是 calculateStatus 中设置 status.state 的连续节选;runtimePods 是刚列出的 Pod,DesiredWorkerReplicas 是期望的 worker 数量:
1 |
|
代入一个 head、两个 worker 的配置,外层条件要求本轮维护没有报错,并且已观察到的 Pod 总数为 3。内层的 CheckAllPodsRunning 再逐个检查这些 Pod,通过后才把状态设为 Ready。仅仅收到了最后一个 worker 的更新通知,还不足以跳过这些检查。
CheckAllPodsRunning 位于 utils/util.go,输入是 Pod 列表,返回值表示是否通过检查。函数如下:
1 |
|
按代码顺序看:列表为空时返回 false;有 Pod 的 phase 不是 Running 时返回 false;有 Ready 条件且值不是 True 时,也返回 false。因此,若三个 Pod 都已 Running,但最后一个 worker 的 Ready 仍为 False,集群就不会在这一分支被设为 Ready。等它变成 True,后续协调重新读取 Pod,满足上述条件后才会设置集群状态。
集群状态计算成功并写回 API 后,RayJob Controller 通过自己注册的 Owns(RayCluster) 观察到变化,把所属 demo-job 的标识放入自己的队列。轮到它处理时,再读取集群状态、判断提交条件。到这里,worker 的就绪变化才传到了作业管理层:两次协调通过 API 中保存的状态接续,RayCluster Controller 完成本轮后即可返回。
前面说明了两层 Controller 怎样通过资源状态交接。每个 Controller 完成本轮能做的事情后,都需要把执行机会交回框架:集群可能仍在准备,Ray 作业也可能还在运行。接下来要解决的是,怎样在条件变化后继续处理这些对象。
一轮 Reconcile 返回以后,watch 仍然持续工作。例如,worker 的 Ready 条件更新,可以沿前面的事件路径触发新的 RayCluster 协调;集群状态写回后,又可以触发 RayJob 协调。每轮处理都会重新读取当前状态、判断下一步条件。这种由事件触发的检查,本身不需要定时轮询。
随着案例进入作业执行阶段,RayJob Controller 还要跟踪另一类进展:Ray 内部的作业状态。程序执行完成时,承载它的 Pod 可能仍然正常运行,未必产生相应的 Kubernetes 资源更新。只等待 Pod 事件,就可能无法及时发现作业已经结束。因此,RayJob Controller 会查询 Ray Jobs API;如果作业还在运行,就安排稍后再查。
能通过监听获知的变化,可以由事件推进。等待集群和跟踪 Ray 作业时,KubeRay 也会安排延时重查;等待期间 watch 继续接收变化,两条路径都可以触发后续协调。
前面写回 API 的 status 描述资源已经走到哪一步,供使用者和其他 Controller 读取。这里的 Result 与 error 则返回给调用 Reconcile 的框架:前者表达是否以及何时再查,后者表达本轮是否失败。框架根据它们安排当前对象的下一轮处理。
例如,成功查询 Ray Jobs API、得知作业仍在运行,是正常等待,可以要求稍后再看;查询失败、没能获得作业状态,则需要按错误规则安排重试。一次协调成功,只表示本轮检查和管理动作没有报错,并不要求 Ray 作业已经完成。
需要延时重查时,代码返回 RequeueAfter,把后续安排交给框架,当前函数随后就结束了。等待期间,协调协程可以处理其他对象,不必留在函数里等待 Pod 就绪或 Ray 作业完成。框架登记同一个对象的延时入队要求,到时再给它处理机会;下一轮会从协调入口重新读取资源,而不是从上次函数退出的位置继续执行。
沿用正在跟踪作业的 RayJob Controller,看它怎样表达这个安排。RayJob.Reconcile 的收尾部分包含下面的代码;两段之间省略了一次指标更新:
1 |
|
这两处都返回了相同的延时,实际是否采用,还要看 error。框架先处理错误:普通 error 非空时走限速重试,并忽略 RequeueAfter。rate limiter 即限速器,在这里决定重试的等待时间,避免失败请求立即反复执行。TerminalError 则是明确标记为“不由本次错误自动重试”的错误类型,后续资源变化仍可触发新的协调。
在 controller-runtime reconcileHandler 中,成功且设置延时的分支使用下面两行连续代码:
1 |
|
Forget(req) 清除这个对象的限速重试记录,AddWithOpts 根据 RequeueAfter 指定的时长安排延时入队,之后由框架再次调用协调函数。清除重试记录不会删除 RayJob,也不会停止监听。
| 返回结果 | 框架如何处理 |
|---|---|
| 普通 error 非空 | 记录错误,按 rate limiter 安排重试;忽略返回值中的延时 |
TerminalError |
本次错误不触发自动限速重试;后续资源事件仍可能再次触发处理 |
error 为空,RequeueAfter > 0
|
清除限速记录,安排延时请求 |
error 为空,RequeueAfter <= 0 且 Requeue: true
|
走限速重排路径 |
| error 为空,空 Result | 本轮不主动安排重查;后续事件仍可以再次入队 |
无论由事件、延时还是错误重试触发,请求都用同一个对象标识回到队列,继续按前面说明的规则合并和处理。例如,RayJob 已经安排延时重查,期间又收到下级 RayCluster 的更新,事件请求就可能让它更早得到检查。延时时间到了,也仍要等可用的协调协程。因此,RequeueAfter 不能理解为精准的固定周期;下一轮总会从读取对象开始。
本例中,worker Pod 的变化通过监听进入 Operator,事件处理器按拥有关系把 RayCluster 标识放进队列。RayCluster Controller 读取当前对象和相关资源,维护 Pod、计算状态并写回 API;RayJob Controller 再通过自己的监听和队列读取这个结果,继续获取 Dashboard 地址、创建 submitter Job。两层 Controller 通过保存下来的资源状态协作。
每一轮只完成当前能做的管理动作。返回以后,监听仍然接收资源变化;需要跟踪 Ray 作业等外部进展时,还可以安排延时查询,失败则按错误路径处理。资源状态交给下一层 Controller,返回值交给运行框架,分别承担业务进度传递和后续处理安排。
队列把同一对象尚待处理的请求合并,协调逻辑依据当前能读到的状态判断;本轮进行期间的新变化,还可以留下下一轮检查。对使用者来说,声明已保存、Pod 已就绪、RayJob 已推进到下一阶段,可以出现在不同时间,它们通过上述过程逐步完成。
下一篇把这条路径放到 100 个 RayCluster 同时变化的场景里,沿入队、请求合并和任务分配的实现,分析不同对象怎样并行,以及同一个对象的多轮检查怎样串行衔接。
2026-09-15 23:30:00
第二篇说明了一个 RayCluster 怎样得到一次检查机会,以及本轮结束后怎样继续处理。现在把范围扩大到同一个 Kubernetes 集群里的 100 个 RayCluster:Operator 怎样在它们之间分配处理机会?配置中的协调并发数设为 2,究竟限制了什么?
下面先看不同 RayCluster 怎样共享协调协程,以及同一个 RayCluster 的多次请求为什么不会并行执行;再进入队列内部,说明请求如何合并、等待和分配。取得请求之后,Controller 还要面对缓存同步延迟和其他写入者,这两类问题决定了资源操作怎样安全接续。
先看 RayCluster Controller。它有自己的工作队列,队列中的请求标识需要检查哪套集群。例如 default/demo-job-abcde,斜杠前是命名空间,后面是对象名;源码中的标准 reconcile.Request 保存这两个字段。队列用这个标识区分对象,后文把它称为“键”。
负责从队列取请求并调用 Reconcile 的执行单元叫协调协程。一轮协调会读取资源、判断差异,完成当前能做的创建、更新等操作,然后返回。协程可以接着处理另一套集群,并不固定属于某个 RayCluster。
KubeRay 的 --reconcile-concurrency 默认值是 1。RayCluster Controller 把它传给框架的 MaxConcurrentReconciles,决定启动多少个这样的协程。下面是 controller-runtime 启动协调协程的代码,仅省略循环内的两行解释性注释:
1 |
|
循环中的 go func() 启动一个 Go 协程;每个协程反复执行 processNextWorkItem,取出请求并完成一轮处理。外面的循环次数就是配置的并发数。wg 用于等待这些协程结束,不参与请求去重。源码里这类执行单元也称 worker,它负责管理资源;Ray worker 则在 Ray 集群里执行用户计算。
假设并发数设为 2,集群 A、B、C、D 都需要检查。其中 A 仍是 default/demo-job-abcde,其余字母代表另外三套 RayCluster。开始时两个协程可以分别处理 A、B;A 的这一轮先结束,空出的协程就可以接手 C。
图中时间向右,条块表示一轮 Reconcile,长度仅作示意。A 的检查结束后,它的 Ray 程序仍然可以继续运行。因而,两个协调协程可以轮流管理很多集群,配置值 2 不代表只能运行两个 RayCluster,也不限制 Ray task、actor 的计算并发。
是否占用这个名额,要看当前函数有没有返回。如果 Controller 正在等待一次 HTTP 请求的响应,这一轮尚未结束,仍占一个名额。如果已经发现 Pod 尚未就绪,返回 RequeueAfter 要求稍后再查,本轮就结束了,延时等待由队列安排,协程可以去处理其他集群。
A 和 B 可以同时处理,同一个 A 的两轮协调则按顺序进行。这样,针对 A 的一次数量判断和资源操作结束后,下一轮才会开始,避免两轮根据各自读到的 Pod 列表同时补建 worker。
这里要分别记录两件事:A 是否还有一次检查等待执行,以及 A 是否已经有一轮正在执行。下面暂时忽略延时和优先级,只看立即请求经过队列整理后的变化。
| 时刻 | 发生的事情 | A 的待处理记录 | A 是否正在处理 |
|---|---|---|---|
| 第一次变化 | 登记一次 A 的检查请求 | 有一份 | 否 |
| 开始协调 | 协程取走这份请求 | 无 | 是 |
| 处理中再次变化 | 为 A 登记下一次检查 | 有一份 | 是 |
| 又来一次变化 | 合并到已有待处理记录 | 仍是一份 | 是 |
| 本轮结束 | 解除 A 的执行锁定 | 保留,等待分配 | 否 |
| 再次分配 | 协程取走下一次请求 | 无 | 是 |
表中间两行描述的是:A 正在处理,同时还有一份 A 等待处理。当前轮可能已经读过 Pod,后面发生的新变化就留给下一轮检查。这份请求会先保留在队列中,等当前轮结束后再交给协调协程。
等待 A 的时候,空闲协程仍可以接手 B、C 等其他对象。同一个待处理请求反复到来时会合并;每一轮实际开始后,再读取当时能看到的资源状态。队列不需要为每条 Pod 通知安排一轮独立调用。
如果 A 处理期间没有留下新的请求,本轮结束后也就没有下一份 A 可以分配。后来再次发生变化,仍可以重新登记。去重的范围是当前的待处理记录。
这也说明了提高并发数的作用范围:如果积压来自许多不同集群,更多协程可以让它们同时得到处理;如果反复变化的主要是 A,A 仍要一轮接一轮执行,其他协程不能把它的两轮协调拆开并行。
前面的表说明了请求在等待和执行之间怎样变化。下面把负责这些动作的协程分开,再看它们如何共享队列记录、交接任务。
Controller 使用 controller-runtime 的优先级队列 priorityqueue。创建入口 NewTypedUnmanaged 默认选择它;KubeRay 沿用这一配置,没有通过 NewQueue 替换队列,也没有关闭 UsePriorityQueue。
继续采用协调并发数为 2 的例子。一个 RayCluster Controller 的请求处理路径涉及以下协程:
| 所在位置 | 分工与数量 | 具体工作 |
|---|---|---|
| 队列内部 | 1 个整理请求的协程 | 处理输入缓冲,合并重复请求,登记延时和优先级 |
| 队列内部 | 1 个管理延时的协程 | 等待延时到期,把对应项转入可处理集合 |
| 队列内部 | 1 个分配任务的协程 | 选择可执行对象,交给正在取任务的协调协程 |
| Controller | 2 个协调协程 | 向队列取任务,取得对象后调用 Reconcile
|
队列创建时分别启动 handleAddBuffer、handleWaitingItems 和 handleReadyItems,对应前三行。最后一行由 Controller 按 MaxConcurrentReconciles 启动。这个例子是三个队列后台协程加两个协调协程,配置中的并发数只控制后者;日志、指标等还有其他辅助协程,所以表格不是进程的协程总数。
前三类协程负责维护请求和分配条件;取得任务后的协调协程负责读取集群资源、创建 Pod、写回状态。它们都在同一进程中工作,队列内部通过锁保护共享记录,通过通知唤醒需要继续工作的协程。
提交请求时,调用方通过 AddWithOpts 交出对象标识和处理选项;选项可以表达延时、优先级或限速重试。入口先把它们存入输入缓冲,再通知后台整理。下面是接收请求的连续节选:
1 |
|
addBuffer 是刚收到的请求暂存区,opts 保存选项,items 保存这次提交的对象标识。追加前后加锁和解锁,是为了让多个提交方安全地访问这个缓冲。最后的 notifyItemAddedToAddBuffer 发出“有新请求”的通知。
此处先追加请求,按对象标识合并的工作留到整理阶段。例如 A 连续提交两次,缓冲里可以同时留着这两次请求。调用 AddWithOpts 返回时,请求已经被接收,随后再参与合并。
输入缓冲暂存的是“检查 A”的请求。后文还会用到对象缓存:它由 informer 维护,保存 RayCluster、Pod 等资源的本地记录,供 Reconcile 读取。两者分别服务于任务安排和资源观察。
整理协程负责把输入缓冲中的请求变成待处理记录。它等待 itemAddedToAddBuffer 通知,运行下面的 handleAddBuffer 循环:
1 |
|
可以按一次通知读这段循环:select 先等待;收到 itemAddedToAddBuffer 通知后,取得队列锁,调用 lockedFlushAddBuffer 整理缓冲,再解锁并回到等待。若收到队列关闭信号 done,这个协程就退出。
整理时,lockedFlushAddBuffer 取走当时已有的一批缓冲记录,交给 lockedAddWithOpts 加入或更新待处理项。一次通知可以带来一批处理,没有固定刷新周期。输入缓冲将接收请求和维护队列结构分开,提交方不必一直持有队列主锁做后续整理。
整理后的待处理项存放在按对象标识索引的 items 中。这里的 items 是整个队列的待处理索引,与前面缓冲记录中“本次提交了哪些对象”的同名字段处在不同层次。
以 A 为例:索引里没有 A,就建立一项;已经有 A,就把新要求合并到现有项。前面表中“仍是一份”对应的就是这一步。索引只保存尚待处理的工作,正在执行的那一轮另有记录。
请求有时需要立即处理,有时要求稍后再查。队列另外维护两个集合,按时间条件组织这些待处理项:
| 集合 | 保存什么 | 能否交给协调协程 |
|---|---|---|
waiting |
尚未到期的延时项 | 先等待时间条件满足 |
ready |
无须延时或已经到期的项 | 还要看对象是否正在处理、是否有空闲协程 |
这里的 ready 表示队列项的时间条件满足,与 Pod 或 RayCluster 的 Ready 状态没有关系。items、waiting 和 ready 是同一批待处理工作的不同组织方式,不是依次执行的三份独立任务。
整理完成后,还要通知谁可以继续工作。下面是 lockedAddWithOpts 收尾的连续节选;两个布尔变量记录这一批处理是否新增 ready 项、是否新增或更新 waiting 项:
1 |
|
新增 ready 项时通知分配协程;新增或更新 waiting 项时通知延时协程。一批请求可能同时产生这两类通知。通知只表示“共享记录发生变化,请重新检查”,具体对象仍保存在队列中。整理协程随后释放队列锁,回到循环等待下一次输入通知,不会在这里执行 Reconcile。
假设 A 要求 10 秒后再查,整理阶段会把它放进 waiting。延时协程 handleWaitingItems 负责等待这个处理时间;即使所有协调协程都在忙,它也可以独立更新队列中的到期状态。
它等待三种信号:队列关闭、新增或更新了延时项,以及先前安排的到期信号。下面保留完整函数。ReadyAt 是项的到期时间,toMove 暂存本轮发现的已到期项,nextReady 是下一次时间信号;Ascend 按 waiting 的顺序遍历,最早到期的项在前面。
1 |
|
可以沿一次唤醒看完整流程:
blockForever 不会发出时间信号,协程先等待延时项的新增或更新通知;队列关闭则直接退出。toMove;遇到第一个尚未到期的项,就按它剩余的等待时间设置 nextReady,停止本轮遍历。若新加入的任务更早到期,这次检查也会调整下一次等待。toMove 中的项从 waiting 删除,清除 ReadyAt,更新排序次序和指标,再加入 ready。代码先收集再移动,把遍历和修改树分成两步。例如,A 到期后会从 waiting 转入 ready;如果 A 的上一轮还没结束,它仍不能立即执行。没有延时的 B 则在整理阶段直接进入 ready,整个过程不需要经过延时协程。
整理请求和交出任务由不同的后台协程负责。要理解交接,先看接收任务的一端:协调协程完成上一轮后,会再次调用队列的取任务方法。Controller 的 processNextWorkItem 调用的是 GetWithPriority,同时取得对象标识和优先级;队列的普通 Get 也会转调这个方法。
如果暂时没有任务,协调协程就在取任务方法中等待。GetWithPriority 先登记一个等待者,再通知任务分配逻辑;下面是连续节选:
1 |
|
waiters 记录有多少协调协程已经来取任务、尚未拿到结果。后面的 select 等待队列关闭,或者从 w.get 通道接收一个任务。这里的等待不会不断轮询,也不会执行 Reconcile。
例如,协调协程 1 调用取任务方法后停在通道接收处;A 的请求整理好以后,分配协程确认有人等待,便可以选择 A,通过通道发送。协调协程 1 收到 A,取任务方法返回,才开始执行 Reconcile(A)。这是同一次交接的两端:协调协程主动申请任务,队列内部通过发送响应它,不会因此额外启动一个协调协程。
分配协程 handleReadyItems 把 ready 中的对象与已经来取任务的协调协程配对。它的外层循环先等待通知,下面是连续节选:
1 |
|
整理协程新增 ready 项、延时协程移入到期项、协调协程来取任务,以及 Done 解除对象锁定,都会通过 readyItemOrWaiterAdded 唤醒它。这些变化都可能让一次新的分配成为可能。
被唤醒后,它先取得队列锁,整理一次输入缓冲,再检查 waiters。如果没有协调协程等待取任务,就结束这次分配尝试、释放锁,回到外层循环等待通知。这里结束的是循环内的一次处理,分配协程本身仍在运行。有人等待时,它才继续遍历 ready,按照优先级和次序选择对象。
选择对象时,需要检查“正在处理”的记录。源码用 locked 保存这些对象的键。下面是跳过已锁定对象、锁定本次选中对象的连续节选:
1 |
|
假设当前候选项是 A。若 A 已在 locked 中,return true 表示继续遍历下一项,不会交出 A,也不会把它的待处理请求丢掉。若 A 没有锁定,则记录这次出队的指标,并将 A 加入 locked。
紧接着是本次交接的连续代码:
1 |
|
这几行减少等待者计数,从待处理索引 items 移除 A,并通过 w.get 通道发送。前面调用取任务方法、正在等待接收的某个协调协程收到 A 后返回,随后开始 Reconcile(A)。toDelete 暂存已经交出的项,这一批遍历结束后再统一从 ready 删除。
如果还有等待者,分配协程继续找下一项;没有等待者就停止遍历。若 ready 中没有可交出的对象,例如剩下的键都已锁定,也会结束本次尝试。最后清理 ready、释放锁,回到外层循环等待新通知。它可以一次交出多个不同对象,但不会在本轮处理完后不停扫描空队列。
此时,待处理的 A 已取走,正在执行的 A 仍记录在 locked 中。若新的 A 到来,便可以重新建立待处理项,同时继续受这把键锁约束。
协调协程通过 defer c.Queue.Done(obj) 保证本轮退出时通知队列。Done(A) 清除 A 的执行锁定并通知任务分配逻辑。之前留下的 A 只有在时间条件满足、协程可用时,才会再次被交出。这样,待处理索引负责合并请求,执行锁定负责避免同一个对象的两轮协调重叠。
回到最初的 A 请求:提交方把它追加到输入缓冲,整理后成为一份待处理记录;有延时就先进入 waiting,到期后进入 ready。协调协程发出取任务请求,分配协程才从 ready 中挑选未锁定的对象并交出;分配路径也会顺手整理尚未处理的输入缓冲。A 的这一轮执行结束,由 Done 解除锁定,期间留下的下一份 A 才有机会继续。
前面分别说明了立即请求和延时请求的去向。两种请求如果指向同一个对象,又会怎样合并?若 A 已经安排延时重查,期间又收到一次立即检查的请求,队列合并时会采用较早的处理时间;优先级有差异时采用较高值。因此,新事件可以让延时项提前进入 ready,但无法绕过 A 当前轮的执行锁定。
在可以分配的 ready 项之间,队列先看优先级,同优先级再看次序。延时项到期进入 ready 时重新取得次序;一项被提升优先级后,排在新优先级已有项之后。整个队列不能视为严格的先进先出(FIFO)。
RequeueAfter 表达的是本轮对后续检查的安排。新事件可能使下一轮提前发生,到期后也可能继续等待空闲协程,因此它不保证精准的固定周期。
队列已经让 A 的两轮协调串行执行,但下一轮能否看到上一轮的写入,还取决于资源读取路径。继续看 A 的 compute worker 组:这个组要求两个同类 worker Pod,当前本地缓存只看到一个。
第一轮向 API Server 请求补建一个 Pod,并收到了成功响应。此时 API Server 中已经有新 Pod,但 Operator 本地缓存可能还没收到对应的监听更新。若第二轮马上读取 Pod 列表,仍可能只看到原来的一个。
| 观察位置 | 创建成功、缓存尚未更新时看到的情况 |
|---|---|
| API Server | 已经保存新 Pod,实际已有两个 |
| Operator 的本地对象缓存 | 仍只记录原来的一个 Pod |
| 刚执行创建的协调逻辑 | 知道创建已经成功,可以把这个结果记下来 |
Manager 为 Controller 提供资源客户端。它的普通 Get、List 可以从对象缓存读取,创建和更新则请求 API Server。写入结果通过监听逐步反映到缓存,因此两条路径之间有同步时间差。启动时会完成缓存和初始事件同步,运行期间的更新则持续到达;即使两轮协调串行执行,后一轮读取时也可能还处在这段同步过程中。
为了让下一轮知道“还有一次创建成功的结果尚未观察到”,KubeRay 保存一份扩缩容预期记录,源码称为 scale expectations。记录包含所属 RayCluster、worker 组、Pod 名称以及创建或删除动作。
登记位置在 createWorkerPod 中:创建请求返回成功之后,才把这个 Pod 记入预期。下面只省略失败分支中的事件记录:
1 |
|
代入本例,这份记录表达的是“已经为 A 的 compute 组成功创建了这个名字的 Pod,等待后续读取确认”。下一轮即使暂时仍只读到一个 Pod,也能先检查是否还有未观察到的创建结果。
处理每个 worker 组之前,Controller 调用 IsSatisfied 检查该组的预期。下面是 reconcilePods 中 worker 组循环开头的判断:
1 |
|
continue 跳过本轮对这个组的 Pod 维护,循环继续检查其他组。协调协程不会停在这里等缓存,下一轮走到这个组时还会重新判断。预期按 namespace、RayCluster 和 group 查询,A 的创建结果暂时不可见,也不会阻止 B 的协调。head 有对应检查,用空字符串作为组标识。
IsSatisfied 先从预期存储中查出本组记录,再逐条调用 isPodScaled 判断。下面保留它的检查循环;items 是查出的记录列表,rp 是当前这条 Pod 创建或删除预期:
1 |
|
遇到未满足的记录就返回 false,已经满足的记录则从预期存储中删除。全部通过,或者已经没有记录时,返回 true,本组继续维护 Pod。清理预期记录表示这次操作已经完成确认,Kubernetes Pod 本身保持原状。
预期随着后续协调中的读取和判断逐步清理。具体条件在 isPodScaled 中:它按记录中的 Pod 名称调用客户端 Get,沿用 Manager 客户端的普通读取路径,因此仍可能从对象缓存获取结果。下面保留完整的两个分支,只省略创建分支中的注释:
1 |
|
创建分支在能读到 Pod 时返回 true,尚不要求它 Ready;读取不成功时,则比较登记时间加上 ExpectationsTimeout 是否早于当前时间。这个值在源码中设为 30 秒,超时也会返回 true,让 IsSatisfied 清掉该记录。
删除分支只有读取返回 NotFound 才返回 true。仍能读到 Pod,或者读取发生其他错误,都会返回 false。Pod 虽然已经标记删除,但只要仍可查询到,删除预期就没有满足。命令行中的 Terminating 表示删除中的显示状态,不是 Pod 的一种 phase。
创建预期保留 30 秒的等待窗口。源码注释给出的例子是:Pod 创建成功后,在下一轮检查前就被其他组件删除,Controller 因而没有机会读到它存在。后续检查发现等待已超时,就清理这条记录,恢复按当前 Pod 列表计算数量差异,决定是否补建。这里恢复的是本组的资源维护,Pod 是否就绪仍由状态检查判断。
等待时长在后续调用 IsSatisfied、检查到这条记录时计算,并没有单独的 30 秒定时器。新的 Pod 事件可以带来下一轮协调,RayCluster 主流程的正常收尾也会返回延时重查要求:ctrl.Result{RequeueAfter: time.Duration(requeueAfterSeconds) * time.Second}。这个分支默认使用 300 秒,可通过环境变量调整;错误或状态变化等路径可能安排更早的处理。因而,30 秒是预期检查使用的超时条件,实际检查时刻由下一轮协调决定。
删除预期以读到 NotFound 作为完成条件,常规检查中没有对应的超时放行。只要记录仍保留、读取尚未返回 NotFound,后续协调就继续暂缓这个组的 Pod 维护,同时处理其他组和其他集群。因此,这个组恢复扩缩容的时机取决于何时确认对象已经消失;等待较久时,也沿用同一条判断规则。
预期也会随所属资源的生命周期清理。例如,RayCluster 读取返回 NotFound 时会清空相关预期;满足条件的 Recreate 升级流程在删除 Pod 后也会清空记录。这些清理由各自的生命周期条件触发,常规删除预期仍按前述读取结果判断,等待时长本身不会触发强制放行或强制清理 Pod。
这些预期保存在 Operator 进程内,底层使用 client-go 提供的线程安全 Indexer,即支持索引查询的内存存储。它保护记录的并发访问,预期本身则帮助下一轮识别缓存尚未反映的操作结果。
expectations 用于减少同一进程前后两轮协调中的重复操作。进程退出后记录随之丢失;API 已创建成功但响应丢失时,也可能还没来得及登记。创建预期超时后,Controller 重新按读到的 Pod 列表判断数量,此时缓存仍可能处于同步过程中。重启后怎样识别已有资源,还要结合资源身份和状态管理,下一篇继续讨论。
到这里,队列控制处理顺序,扩缩容预期补充缓存尚未反映的操作结果。它们都属于 Controller 本地的管理机制。扩大到不同 Controller、不同进程和使用者同时操作时,还要区分限制的范围。
RayJob Controller 可以查询作业状态,同时由 RayCluster Controller 检查 worker Pod。二者在 main.go 中分别注册,各有队列和协调协程,共享 Manager 提供的客户端等基础设施。
同一个 config.ReconcileConcurrency 分别传给 RayCluster、RayJob、RayService。设为 2 时,每个 Controller 各有两个协调名额。如果三边都有足够的待处理请求,合计可以有六轮协调同时进行。
图中的队列位于各自的 Controller 内。队列键包含 namespace,所以 team-a/demo-job-abcde 与 team-b/demo-job-abcde 是两个不同对象,可以并发协调。这只是同名对象的对比例子,主案例 A 仍在 default 命名空间。namespace 用来区分资源,不会自动分得独立队列、专属协程或公平份额。
KubeRay 的 --watch-namespace 可以限定一个或逗号分隔的多个 namespace,空值表示监听所有 namespace。它决定缓存观察的范围;Operator 身份能否读取和修改这些资源,还由 Kubernetes 的访问权限规则 RBAC 决定。
100 个 RayCluster 可以都在同一个 Kubernetes 集群中。这里的启动入口用一份 rest.Config 创建一个 Manager,这份客户端配置指定一套 Kubernetes API 地址和访问身份。增加 namespace 不会增加连接的 Kubernetes 集群;管理另一套 Kubernetes 集群,需要为它配置对应的 Operator 管理实例。
增加 Operator 副本还涉及选主:多个实例中由一个主实例负责协调,其余等待接管。KubeRay 默认开启 leader election,所以共享同一选举锁的副本,不能简单把协调并发数相加。若独立实例的观察范围重叠,各自队列的键锁也不能互相排斥;选举和接管在第五篇展开。
队列的串行规则作用于本 Controller 对同一个 RayCluster 的协调。与此同时,用户仍可以修改对象,其他 Controller 也可以写入相关资源。不同队列或不同对象的协调若涉及同一个下级资源,就需要由资源更新时的并发检查处理;进程内的共享内存则由相应的同步机制保护。
以一次状态更新为例:Controller 读到 A 的配置和状态,准备计算新的 status;在它提交之前,用户修改了 A 的 spec,API Server 已经保存了新版本。若直接用刚才读到的旧对象写回,就需要判断这个写入是否还基于有效版本。
Kubernetes 用 resourceVersion 标记对象版本。读取时客户端拿到这个标记,更新时原样带回,不自行推算或递增。普通带版本的 Update 若携带旧版本,API Server 会拒绝更新并返回 Conflict。这样,读取期间允许其他人修改,提交时再检查版本,这种方式称为乐观并发检查。
RayCluster 的状态写回函数 updateRayClusterStatus 使用 r.Status().Update(ctx, newInstance),把失败交回主流程。主流程在选择错误后执行下面的返回判断:
1 |
|
这段代码可以同时返回错误和等待时间,实际重试规则由框架的 reconcileHandler 决定。普通错误非空时,通过 AddWithOpts 的 RateLimited 选项安排限速重试,并忽略 RequeueAfter;终止型错误有独立分支,不由本次错误自动重试。没有错误时,正的 RequeueAfter 才用于安排延时检查。
这次版本冲突后,本轮结束,后续协调重新读取、计算,再尝试写回。重新读取是其中必要的一步:同一份旧对象仍携带旧版本,原样重试还会冲突;下一轮能否写入,也取决于读到的版本是否已经更新。
版本检查以单次对象更新为单位。假设本轮先创建了一个 Pod,随后写 RayCluster status 时发生冲突,已创建的 Pod 会保留,下一轮结合配置、实际 Pod 和预期记录继续判断。这些跨对象操作分别生效,因此资源状态和 status 需要在后续协调中继续对齐。API 的 Patch 请求只描述部分修改,其冲突条件取决于具体写法,与这里的 Update 规则需要分别理解。
回到 100 个 RayCluster:全部稳定运行,与同时扩容、同时创建大量 Pod,给 Operator 带来的负载不同。稳定集群仍有周期性检查:RayCluster 主流程在收尾时安排下一轮,间隔可由环境变量配置;检查发现资源符合期望,就可以直接结束。
判断并发数是否合适,可以沿请求的处理过程看:等待处理的对象是否积压,协调协程是否长期占满,一轮处理是不是变慢,以及错误重试是否增加。controller-runtime 提供的指标分别对应这些位置:
| 要观察的情况 | 对应指标 |
|---|---|
| 已开始、尚未返回的协调是否长期占满名额 |
controller_runtime_active_workers、controller_runtime_max_concurrent_reconciles
|
| 待处理项有多少,取得处理机会前等了多久 |
workqueue_depth、workqueue_queue_duration_seconds
|
| 一轮协调本身耗时多久 | controller_runtime_reconcile_time_seconds |
| 协调失败是否增多,是否带来更多重试 | controller_runtime_reconcile_errors_total |
这些指标有统计范围。priorityqueue 从项进入 ready 集合开始计量:尚未到期的延时项不计入 depth,排队耗时也不包含前面的延时等待。depth 带有 priority 标签,查看一个 Controller 的总积压时需要汇总各优先级。它表示队列积压,不能作为 RayCluster 总数或尚未 Ready 的集群数。
如果协调协程长期占满、许多不同对象持续排队,而单轮耗时和 API 请求表现稳定,提高并发值得验证。若积压主要来自一个正在处理的对象,仍受同键串行限制;若 Kubernetes API 已经变慢或限流等待增加,继续加协程还可能增加争用。
KubeRay 注册的 rest_client_request_duration_seconds 和 rest_client_rate_limiter_duration_seconds 分别帮助观察 API 请求耗时和客户端限流等待。客户端的 QPS、Burst 在 Manager 创建前配置,分别约束持续请求速率与突发额度,判断并发效果时也要结合这些限制。
协调协程主动取任务,队列后台分别整理请求、处理延时和分配任务;取得对象后,协调协程才执行 Reconcile。同一个集群的待处理请求可以合并,本轮执行期间也可以留下下一轮请求。待处理索引和执行锁定共同保证这些后续检查能够保留,同时避免同一对象的两轮协调重叠。
进入资源维护后,扩缩容预期衔接 API 写入与缓存观察之间的时间差,对象版本检查处理多个写入者的更新冲突。队列、预期记录和版本检查分别作用于请求安排、操作确认和资源写回,合在一起说明了多套集群的管理工作怎样推进。
评估并发配置时,需要结合每个 RayCluster 的 Pod 规模和变更负载,比较调整前后的排队耗时、错误速率及集群就绪耗时,才能判断这 100 个集群是否得到及时处理。下一篇把视角转向进程重启:资源已经创建,状态尚未写回时,Operator 怎样识别已有进展并继续协调。
2026-09-15 23:25:00
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 不代建。主案例采用的是新建专用集群。
集群入口准备好后,同一个 Initializing 分支进入提交用 Job 的准备。createK8sJobIfNeed 按 RayJob 的命名空间和名称查询 Kubernetes Job:主案例中它也叫 demo-job。找到就返回成功,确认不存在才构造并创建。主流程随后执行下面的代码:
1 |
|
这里先确保 Job 存在,再把管理状态设为 Running,最后才通过统一写回保存。创建 Kubernetes Job 和更新 RayJob 是两个 API 请求,图 3 展示其中的中断位置。
可以沿三个时间点看恢复过程:
Initializing,下一轮查询不到同名 Job,继续创建。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 命令是否保留预查询,取决于其实现。换了集群、标识或原记录已删除、丢失时,需要按新的记录状态处理。
新 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 与持续服务入口,继续区分每一层保存的状态和能够恢复的范围。