KubeRay 源码阅读:从 RayJob 看组件分工与协作
第零篇从计算与部署需求解释了 Ray、Kubernetes 和 KubeRay 的分工。这篇沿一次 RayJob 提交进入源码,追踪资源之间的交接。
向 Kubernetes 提交一个 RayJob 后,集群里会陆续出现 RayCluster、head Pod、worker Pod,以及一个用来提交作业的 Kubernetes Job。用户只创建了一个对象,剩下的资源是谁创建的?这些资源准备好之后,程序又由谁启动?
本例使用 default 命名空间中的 RayJob demo-job:新建专用集群,包含一个 head Pod 和 compute 组的两个 worker Pod,采用 K8sJobMode,也就是由 Kubernetes Job 运行提交程序的模式,成功完成后等待 60 秒清理集群。程序及依赖已经放进运行镜像,暂不启用自动扩缩容。
worker group 是一组采用相同模板和资源配置的 worker,compute 是本例给这组起的名字。后文谈到组的副本数,指希望创建多少个相应的 worker Pod;更复杂的多 host 配置留到第九篇说明。
下文把这个提交用 Job 称为 submitter Job,它创建的 Pod 称为 submitter Pod。第一至第五篇的基础案例不启用 Volcano;第六至第九篇再增加批调度配置,分析多个作业竞争资源时的行为。
阅读时会反复遇到一个问题:当前这一步做完以后,交给谁继续?把这些交接点找出来,作业停在“集群已创建,程序还没运行”时,就有了排查入口。队列实现、故障恢复和多副本选主在后续文章中展开。
本文基于 KubeRay 提交 6bf05eb17a3e,涉及 Ray 进程结构时参考 Ray 提交 3f785c0711b9。这是对固定源码版本的分析,文中的时序图表示代码路径,没有加入集群实测结果。
1. 谁管理集群,谁执行计算
沿用导引中的分工: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 运行时和提交程序分别运行在其他容器里,后面会展开它们的位置。
2. RayJob 为什么还要创建 RayCluster
使用 KubeRay 时,最常见的入口是 RayCluster、RayJob 和 RayService。
CRD 定义自定义资源的结构,CR 是它的具体实例。安装 RayJob CRD 以后,就可以创建 demo-job 这样的 RayJob 对象。Kubernetes 保存对象,Controller 读取声明并执行管理逻辑。
选择哪种资源,取决于希望 KubeRay 帮忙管理到哪一步。
| 资源 | 使用者声明什么 | KubeRay 主要管理什么 |
|---|---|---|
| RayCluster | 一套 Ray 集群的配置,包括 head 和各 worker group | 集群所需的 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 使用已有集群,此时不创建新的专用集群,收尾时也要遵循已有集群的生命周期约束。
3. Controller 和 Ray 进程分别运行在哪里
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 是两回事。
4. Controller 怎样把声明变成运行状态
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。
5. 跟踪一次 RayJob 提交
下面摘出 demo-job 中与提交、清理有关的 YAML。必填的集群和 Pod 模板已省略,片段不能直接提交;Ray 运行环境中已经准备好 /app/main.py 及其依赖。
1 | |
这几种对象的名称会在后文反复出现,先把它们对应起来。表中的随机后缀是示意值,实际值应从 RayJob 的 status 读取。
| 对象或标识 | 本例中的值 | 从哪里来 |
|---|---|---|
| RayJob | default/demo-job |
使用者在 YAML 中指定 |
| 专用 RayCluster | default/demo-job-abcde |
Controller 自动生成,保存到 status.rayClusterName |
| worker group | compute |
rayClusterSpec 中的组名,本例要求两个 worker Pod |
| 提交用 Kubernetes Job(submitter Job) | default/demo-job |
沿用 RayJob 的名称;资源类型不同,所以可以同名 |
| Ray 作业提交标识(submission ID) | demo-job-fghij |
本例未设置 spec.jobId,由 Controller 生成并保存到 status.jobId |
集群名和 submission ID 各自生成随机后缀,不能靠一个推导另一个。status.jobId 对应 Ray CLI 的 --submission-id,也不是 Ray 内部给 driver 分配的 JobID。submitter Job 创建的 Pod 另有名称,不能把 Job 名称当作 Pod 名称使用。第四篇讨论整次执行重试时,还会看到集群名和 submission ID 重新生成。
本例已在镜像里准备好程序。创建 RayJob 不会自动上传提交者电脑上的目录;采用挂载文件或 Ray runtime_env 分发代码时,也需要另外配置。runtime_env 是 Ray 描述作业运行环境的机制,可以指定代码、Python 依赖等内容。
5.1 准备集群
图中省略事件传递和重复检查,只保留准备集群的关键步骤。“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 提交。
5.2 提交程序并跟踪执行
下面接着看提交、执行和清理。图中的 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 启动说明提交程序开始运行,二者都还没走到用户程序完成这一步。
5.3 结束后保留什么
本例设置了 shutdownAfterJobFinishes: true 和 ttlSecondsAfterFinished: 60,使用这一收尾路径,且没有启用删除 RayJob CR 的环境变量。对于正常完成的作业,Controller 按结束时间计算等待期限,到期后删除专用 RayCluster,再由 Kubernetes 按资源拥有关系回收下级资源。
handleShutdownAfterJobFinishes 在这条路径中保留 RayJob CR 和提交用的 Kubernetes Job。这样还能查询 RayJob 的状态,并在相关 Pod、日志仍可访问时查看提交日志。
结果文件、日志和程序版本仍需另外归档。保留下来的 RayJob 对象主要记录这次执行的声明与状态。
6. 从资源状态找到负责的组件
如果集群已经创建、程序却还没运行,可以先看流程停在哪两个资源之间,再去查负责这一段的组件。
| 观察到的情况 | 沿本文路径,下一步检查什么 |
|---|---|
| RayJob 已存在,还没有专用 RayCluster | RayJob 状态、事件,以及 RayJob Controller 创建集群的分支 |
| RayCluster 已存在,Pod 尚未创建 | RayCluster Controller 的资源创建过程及相关错误 |
| Pod 已创建,但一直 Pending | 先看是否已分配 Node;未分配时查调度条件,已分配时查镜像拉取、存储挂载等启动步骤 |
| 集群已经 Ready,还没有 submitter Job | RayJob 的准备状态、Dashboard 地址获取和提交用 Job 的创建过程 |
| submitter 已运行,Ray 中找不到作业 | submitter 日志、到 Jobs API 的连接及提交响应 |
| Ray 作业已完成,RayJob 状态尚未完成 | Operator 的状态查询,以及 submitter Job 是否已经结束 |
这些是检查入口,具体原因还要结合事件和日志判断。同一套分工也能解释另外两种用法:直接使用 RayCluster 时,使用者自行安排作业提交;使用 RayService 时,上层 Controller 改为通过 Serve API 管理应用、检查健康并维护入口。需要替换集群的更新路径还涉及 active(当前服务集群)和 pending(准备接替的集群)的切换,对应实现位于 rayservice_controller.go,第五篇会展开。
下一篇继续看资源变化怎样触发协调:一个 worker Pod 的事件,如何映射到对应的 RayCluster,再进入 Controller 的工作队列。
参考资料
本文 Go 节选仅调整缩进;拼接和省略范围已在上下文标明。
- KubeRay 源码基线(6bf05eb17a3e)
- Ray 源码基线(3f785c0711b9)
- KubeRay:ray-operator/controllers/ray/rayjob_controller.go
- KubeRay:ray-operator/main.go
- Ray:python/ray/_private/node.py
- KubeRay:ray-operator/controllers/ray/common/pod.go
- Ray:python/ray/autoscaler/_private/kuberay/node_provider.py
- KubeRay:ray-operator/controllers/ray/raycluster_controller.go
- KubeRay:ray-operator/controllers/ray/common/job.go
- KubeRay:ray-operator/controllers/ray/utils/dashboardclient/dashboard_httpclient.go
- KubeRay:ray-operator/controllers/ray/rayservice_controller.go
- KubeRay:ray-operator/apis/ray/v1
- KubeRay:ray-operator/controllers/ray/utils/util.go
- KubeRay:ray-operator/controllers/ray/common/association.go
- Ray GCS 概念与容错条件
- Ray Core 概念与环境依赖
- Kubernetes:annotations