第 2 篇|KubeRay 协调:从 Pod 变化到 Reconcile

本文最后更新于:2 天前

前言

第一篇里,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 中判断。

图 1:Pod 状态变化进入 Controller 的路径
图 1:Pod 状态变化进入 Controller 的路径 查看原图

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
2
3
4
5
6
7
8
9
10
b := ctrl.NewControllerManagedBy(mgr).
For(&rayv1.RayCluster{}, builder.WithPredicates(predicate.Or(
predicate.GenerationChangedPredicate{},
predicate.LabelChangedPredicate{},
predicate.AnnotationChangedPredicate{},
))).
Owns(&corev1.Pod{}).
Owns(&corev1.Service{}).
Owns(&corev1.Secret{}).
Owns(&corev1.PersistentVolumeClaim{})

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
2
3
request := reconcile.Request{NamespacedName: types.NamespacedName{
Name: ref.Name,
}}

后续代码根据资源作用域补上 namespace。对 RayCluster 来说,请求标识形如 default/demo-job-abcde,只告诉协调函数该检查哪套集群。至于要不要创建 worker,要等读到对象后再判断。

前面的 For 和 Owns 接收不同来源的事件,最后都可以得到同一个 RayCluster 标识,两条入口由此汇合到同一套协调逻辑。RayJob 对下级 RayCluster 也有自己的监听,形成图中的两层关系。图里的每一层都先把所属对象入队,再由对应 Controller 处理;实际的状态交接要等协调函数执行以后才会发生。

图 2:Pod、RayCluster 和 RayJob 之间的事件接力
图 2:Pod、RayCluster 和 RayJob 之间的事件接力 查看原图

待处理请求的合并

假设 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
2
3
4
instance := &rayv1.RayCluster{}
if err = r.Get(ctx, request.NamespacedName, instance); err == nil {
return r.rayClusterReconcile(ctx, instance)
}

这段代码用请求里的 namespace/name 读取 RayCluster,把结果放进 instance,读取成功才进入实际维护函数。后面的 NotFound 分支处理对象已删除、旧请求仍留在队列里的情况,此时可以结束本轮;连接或权限等读取错误则要按失败处理。

进入 rayClusterReconcile 后,代码会检查资源是否交给外部 Controller 管理,验证配置,处理删除流程,再按当前状态维护集群资源。不同配置会走不同分支。本例没有删除请求,也不启用自动扩缩容,重点是 head Service、head Pod、worker Pod 和集群状态。

图 3:一轮 RayCluster 协调做什么
图 3:一轮 RayCluster 协调做什么 查看原图

spec 保存使用者的期望,status 保存 Controller 已观察到的集群情况;实际 Pod、Service 等对象也参与判断。只读某一个 status 字段,不能代替完整的资源检查。

假设配置要求两个 worker,本轮只观察到一个,Controller 会结合已有 Pod 和尚未被缓存观察到的创建记录,判断是否需要补建。创建请求发出后,Pod 可能还要等待调度和镜像拉取,本轮只完成当前能做的管理动作;如果数量已满足,就继续检查这些 Pod 的状态。

计算集群状态与 Controller 交接

资源维护之后,还要把观察结果整理成集群状态。在 raycluster_controller.go 中,rayClusterReconcile 执行完各项维护操作,接着调用负责计算状态的 calculateStatus:

1
newInstance, calculateErr := r.calculateStatus(ctx, instance, reconcileErr)

这里的 instance 是当前 RayCluster,reconcileErr 记录前面资源维护是否出错。calculateStatus 复制这个对象,列出所属 Pod,再把计算结果放进返回的 newInstance;调用方在计算成功后将新状态写回 API。计算状态和写回状态是这一步里的两项操作。

本例关注集群第一次进入 Ready 的条件。下面是 calculateStatus 中设置 status.state 的连续节选;runtimePods 是刚列出的 Pod,DesiredWorkerReplicas 是期望的 worker 数量:

1
2
3
4
5
6
if reconcileErr == nil && len(runtimePods.Items) == int(newInstance.Status.DesiredWorkerReplicas)+1 { // workers + 1 head
if utils.CheckAllPodsRunning(ctx, runtimePods) {
newInstance.Status.State = rayv1.Ready
newInstance.Status.Reason = ""
}
}

代入一个 head、两个 worker 的配置,外层条件要求本轮维护没有报错,并且已观察到的 Pod 总数为 3。内层的 CheckAllPodsRunning 再逐个检查这些 Pod,通过后才把状态设为 Ready。仅仅收到了最后一个 worker 的更新通知,还不足以跳过这些检查。

CheckAllPodsRunning 位于 utils/util.go,输入是 Pod 列表,返回值表示是否通过检查。函数如下:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
func CheckAllPodsRunning(ctx context.Context, runningPods corev1.PodList) bool {
log := ctrl.LoggerFrom(ctx)
// check if there are no pods.
if len(runningPods.Items) == 0 {
return false
}
for _, pod := range runningPods.Items {
if pod.Status.Phase != corev1.PodRunning {
log.Info("CheckAllPodsRunning: Pod is not running.", "podName", pod.Name, "pod Status.Phase", pod.Status.Phase)
return false
}
for _, cond := range pod.Status.Conditions {
if cond.Type == corev1.PodReady && cond.Status != corev1.ConditionTrue {
log.Info("CheckAllPodsRunning: Pod is not ready.", "podName", pod.Name, "pod Status.Conditions[PodReady]", cond)
return false
}
}
}
return true
}

按代码顺序看:列表为空时返回 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
2
3
4
5
6
if err = r.updateRayJobStatus(ctx, originalRayJobInstance, rayJobInstance); err != nil {
logger.Info("Failed to update RayJob status", "error", err)
return ctrl.Result{RequeueAfter: RayJobDefaultRequeueDuration}, err
}
// ... 指标更新省略。
return ctrl.Result{RequeueAfter: RayJobDefaultRequeueDuration}, nil

这两处都返回了相同的延时,实际是否采用,还要看 error。框架先处理错误:普通 error 非空时走限速重试,并忽略 RequeueAfter。rate limiter 即限速器,在这里决定重试的等待时间,避免失败请求立即反复执行。TerminalError 则是明确标记为“不由本次错误自动重试”的错误类型,后续资源变化仍可触发新的协调。

在 controller-runtime reconcileHandler 中,成功且设置延时的分支使用下面两行连续代码:

1
2
c.Queue.Forget(req)
c.Queue.AddWithOpts(priorityqueue.AddOpts{After: result.RequeueAfter, Priority: new(priority)}, req)

Forget(req) 清除这个对象的限速重试记录,AddWithOpts 根据 RequeueAfter 指定的时长安排延时入队,之后由框架再次调用协调函数。清除重试记录不会删除 RayJob,也不会停止监听。

图 4:Reconcile 返回值如何影响后续调度
图 4:Reconcile 返回值如何影响后续调度 查看原图
返回结果 框架如何处理
普通 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 同时变化的场景里,沿入队、请求合并和任务分配的实现,分析不同对象怎样并行,以及同一个对象的多轮检查怎样串行衔接。

参考资料


第 2 篇|KubeRay 协调:从 Pod 变化到 Reconcile
https://tanxinyu.work/kuberay-volcano-02-reconcile/
作者
谭新宇
发布于
2026年9月15日
许可协议