第 1 篇|KubeRay 架构:组件分工与协作

本文最后更新于:2 天前

前言

第零篇解释了 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 说明接入方式。本文先走完资源创建、程序提交和收尾,队列、恢复与多副本接管留到后文。

KubeRay 的职责与核心资源

集群管理与计算执行

沿用第 0 篇中的分工:Kubernetes 将 head 和 worker Pod 部署到节点上;Ray 在容器内安排 task 和 actor;KubeRay Operator 则根据 RayJob 等声明准备资源、跟踪状态。下面把这层分工落到两个具体接口上:Kubernetes API 管理资源对象,Ray Jobs API 接收入口命令并提供作业状态查询。

图 1:KubeRay 的管理路径与 Ray 的计算路径
图 1:KubeRay 的管理路径与 Ray 的计算路径 查看原图

顺着图中的连接看,Operator 一侧通过 Kubernetes API 管理资源,另一侧通过 Ray 管理 API 查询作业或管理 Serve 应用。Ray 节点之间传输对象、执行 task 和 actor,不需要经过 Operator。

Operator 的主体是 Go 程序。Ray 运行时和提交程序分别运行在其他容器里,后面会展开它们的位置。

RayJob、RayCluster 与 RayService

使用 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 的集群管理逻辑。

图 2:三种资源对应的管理关系
图 2:三种资源对应的管理关系 查看原图

图中的三列是独立的使用方式,分别对应不同的 RayCluster;连线表示资源的创建和管理关系,辅助资源有所省略。下文沿中间一列的新建专用集群路径阅读。

以 RayJob 新建专用集群为例,它的 spec.rayClusterSpec 中包含集群配置。RayJob Controller 会根据这份配置构造一个独立的 RayCluster 对象。下面是 constructRayClusterForRayJob 中的连续节选:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
rayCluster := &rayv1.RayCluster{
ObjectMeta: metav1.ObjectMeta{
Labels: labels,
Annotations: annotations,
Name: rayClusterName,
Namespace: rayJobInstance.Namespace,
},
Spec: *rayJobInstance.Spec.RayClusterSpec.DeepCopy(),
}

// Set the ownership in order to do the garbage collection by k8s.
if err := ctrl.SetControllerReference(rayJobInstance, rayCluster, r.Scheme); err != nil {
return nil, err
}

DeepCopy() 把 RayJob 内嵌的集群配置复制给新对象。Labels 和 Annotations 都是对象上的键值元数据:label 常用于筛选资源,annotation 常用于附加说明或组件约定的配置,具体含义由读取它的程序解释。随后设置的 owner reference 记录 RayCluster 属于这个 RayJob,供 Kubernetes 垃圾回收等机制使用;这里的垃圾回收指依据资源拥有关系清理下级对象。新对象提交到 API 后,具体的集群创建工作就由 RayCluster Controller 接手。

RayJob 也支持通过 clusterSelector 使用已有集群,此时不创建新的专用集群,收尾时也要遵循已有集群的生命周期约束。

组件部署与协作方式

Controller 与 Ray 进程

KubeRay 核心 Operator 的一个实例运行一个 Go 主进程。controller-runtime 是用于编写 Kubernetes 控制器的 Go 库,其中的 Manager 对象负责组织 Controller 的运行,共享客户端、缓存等基础设施。RayCluster、RayJob 和 RayService 各有自己的 Controller,但它们注册在同一个 Manager 上。

main.go 中的三处注册如下。这里拼接了三段原始调用,省略中间的选项构造和其他可选组件注册。

1
2
3
4
5
6
7
8
9
10
exitOnError(ray.NewReconciler(mgr, rayClusterOptions).SetupWithManager(mgr, config.ReconcileConcurrency),
"unable to create controller", "controller", "RayCluster")

// ...
exitOnError(ray.NewRayServiceReconciler(ctx, mgr, config).SetupWithManager(mgr, config.ReconcileConcurrency),
"unable to create controller", "controller", "RayService")

// ...
exitOnError(ray.NewRayJobReconciler(ctx, mgr, rayJobOptions, config).SetupWithManager(mgr, config.ReconcileConcurrency),
"unable to create controller", "controller", "RayJob")

这三个调用传入的是同一个 mgr,所以三个 Controller 都运行在同一个 Operator 进程里。增加 Operator 副本时,增加的是整套进程实例。

图 3:Operator 与 Ray 的部署结构
图 3:Operator 与 Ray 的部署结构 查看原图

图 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
2
3
4
5
6
// SetupWithManager 中的链式调用开头,后续 WithOptions 和 Complete 省略。
return ctrl.NewControllerManagedBy(mgr).
For(&rayv1.RayJob{}).
Owns(&rayv1.RayCluster{}).
Owns(&corev1.Service{}).
Owns(&batchv1.Job{}).

For 指定主要处理的资源类型。Owns 声明对相关从属资源的关注,这些资源发生事件时,可以触发其所有者的协调。事件如何映射回某个 RayJob、如何进入队列,下一篇再顺着这段代码往下读。

RayJob Controller 与 RayCluster Controller 通过 Kubernetes 资源及其状态交接工作:前者创建集群声明,后者创建 Pod、观察就绪情况并更新集群状态,前者再读取这个结果。整个过程跨越多轮 Reconcile。

一次 RayJob 的完整生命周期

示例配置与名称说明

把前面的分工放进一个具体例子:运行镜像里已经准备好 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
2
3
4
5
6
7
8
9
10
11
apiVersion: ray.io/v1
kind: RayJob
metadata:
name: demo-job
namespace: default
spec:
entrypoint: python /app/main.py
submissionMode: K8sJobMode
shutdownAfterJobFinishes: true
ttlSecondsAfterFinished: 60
# rayClusterSpec 在此省略:配置一个 head、compute 组的两个 worker 及其 Pod 模板。

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 依赖等内容。

准备集群

图 4:从 RayJob 到就绪的 RayCluster
图 4:从 RayJob 到就绪的 RayCluster 查看原图

图中省略事件传递和重复检查,只保留准备集群的关键步骤。“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
2
3
4
5
6
7
8
if clientURL := rayJobInstance.Status.DashboardURL; clientURL == "" {
if rayClusterInstance.Status.State != rayv1.Ready {
// ... 原始日志与解释性注释省略。
rayJobInstance.Status.RayClusterStatus = rayClusterInstance.Status
break
}
// ... 查询 head Service 的 Dashboard 地址,并保存到 status。
}

集群尚未就绪时,break 退出当前状态分支,后续代码更新状态并安排再次检查。等到观察到集群就绪,Controller 才会获取 Dashboard 地址并推进提交。这里的 Go 常量 rayv1.Ready 对应 RayCluster 的 status.state: ready,与 Pod 的 Ready 条件、RayJob 的 Running 阶段分别属于不同资源和字段。到这一步,准备好的是运行环境,入口程序还要通过 Jobs API 提交。

提交程序与跟踪执行

下面接着看提交、执行和清理。图中的 TTL 指完成后保留资源的等待期限,本例设为 60 秒。

图 5:提交 Ray 作业、回写状态与清理集群
图 5:提交 Ray 作业、回写状态与清理集群 查看原图

集群准备好后,沿图 5 继续看一次成功的提交和收尾。图中的 Kubernetes 一列合并了 Job 控制器、调度和节点执行环节。

RayJob Controller 在 K8sJobMode 分支中创建提交用的 Kubernetes Job。下面是对应代码的连续节选:

1
2
3
4
5
6
7
if rayJobInstance.Spec.SubmissionMode == rayv1.K8sJobMode {
if err := r.createK8sJobIfNeed(ctx, rayJobInstance, rayClusterInstance); err != nil {
return ctrl.Result{RequeueAfter: RayJobDefaultRequeueDuration}, err
}
}

rayJobInstance.Status.JobDeploymentStatus = rayv1.JobDeploymentStatusRunning

Kubernetes 自己的 Job 控制器根据这个 Job 创建 submitter Pod。submitter 运行 Ray CLI,通过 head Service 访问 Dashboard 的 Jobs API。

BuildJobSubmitCommand 负责拼接提交命令。把健康检查、已有 job 检查和日志跟踪等外围逻辑拿掉,其核心命令形态如下。这是带占位符的示意,不是完整生成结果:

1
2
3
4
5
ray job submit \
--address=http://<head-service>:8265 \
--submission-id=<submission-id> \
--no-wait \
-- python /app/main.py

这条提交命令带有 --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 对象主要记录这次执行的声明与状态。

小结:KubeRay 的分层设计

一次 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 节选仅调整缩进;拼接和省略范围已在上下文标明。


第 1 篇|KubeRay 架构:组件分工与协作
https://tanxinyu.work/kuberay-volcano-01-architecture/
作者
谭新宇
发布于
2026年9月15日
许可协议