KubeRay 源码阅读:一次 Pod 变化怎样触发 Reconcile
第一篇里,RayJob default/demo-job 创建了专用 RayCluster,示例名称为 default/demo-job-abcde,包含一个 head Pod 和 compute 组的两个 worker Pod。RayJob Controller 要等集群就绪,才会创建 submitter Job。假设最后一个 worker Pod 已经进入 Running,Ready 条件也刚变为 True,这个变化怎样传到 Operator,又怎样让作业继续提交?
这篇从 Pod 的变化往下追:事件从哪里进入 Controller,队列保存了什么,一轮 Reconcile 返回后又由谁安排下一轮。把这条链路串起来,才能解释“资源已经变了,Operator 却没有继续处理”可能卡在哪里。
阅读基线是 KubeRay 6bf05eb17a3e 和它依赖的 controller-runtime v0.24.1。图中展示源码路径,时间间隔为示意,没有运行集群实验。
1. Operator 怎么知道 Pod 变了
导引中负责回报 Pod 状态的 kubelet,会把本节点观察到的结果写入 API Server。Operator 通过 watch 接收所关注资源的变化通知;这是 Kubernetes API 提供的监听机制。
管理这些对象同步和通知的客户端组件叫 informer:它维护本地缓存,并把变化交给注册的事件处理器,源码中通常称为 handler。处理器再决定把哪个对象放进待协调队列。第一篇所说的 Controller 通过资源状态交接,在这里展开成了这条通知链路。
这里要分清缓存和队列里各存了什么。本地缓存保留资源对象,供 Controller 读取;队列保存待处理对象的标识。在处理 Pod 事件时,Controller 通常不需要把整个 Pod 快照塞进队列,稍后进入 Reconcile,会重新读取所需资源。
KubeRay 在 main.go 中创建 Manager,配置缓存、监听的 namespace 和客户端。各个 Controller 使用这个 Manager 注册监听,所以一个 Operator 实例中的几个 Controller 可以共享基础设施。
这里的“监听”有范围。WatchNamespace 决定关注哪些 namespace;KubeRay 还在 internal/managercache/cache.go 中为 Pod、Kubernetes Job 等对象设置标签筛选。普通业务 Pod 的每次更新没有必要都送到 KubeRay。监听范围之外还要看权限:RBAC 是 Kubernetes 的基于角色的访问控制,用来决定这个 Operator 身份能对哪些资源执行哪些操作。配置了监听范围,不等于已经获得读取权限。
Controller 开始处理对象之前,需要先完成缓存的初始同步。controller-runtime 的 Kind.Start 注册事件处理器后,会等待缓存同步,并等待这个处理器收到初始事件。Controller 启动协调协程之前,还要等待这些事件源同步完成。goroutine 是由 Go 运行时调度的执行单元,多个 goroutine 可以在同一进程里共享内存;这里把从队列取对象、调用 Reconcile 的 goroutine 称为协调协程,与执行用户计算的 Ray worker 区分。因此,启动时 watch 权限不足、资源类型不可用或缓存同步失败,都可能让 Controller 无法进入正常处理阶段。
“写 API 成功”和“缓存里已看到新对象”之间存在时间差。Manager 提供的常规类型客户端通常从缓存读,写操作直接发给 API Server;绕过缓存的读取另有路径。这个差异在第三篇分析并发创建时还会用到。
2. 一个 Pod 事件为什么会处理 RayCluster
RayCluster Controller 注册监听的代码位于 SetupWithManager。其中 predicate 是事件筛选条件,generation 是用于识别期望配置变更的计数。下面保留构建器调用,省略后面的调度器配置和 Controller 选项:
1 | |
For 指定这个 Controller 主要协调 RayCluster。RayCluster 自身发生符合条件的事件时,处理器把它的 namespace 和 name 放进队列。
这段注册还包括 Secret 和 PersistentVolumeClaim(PVC)。前者用于存放密码等敏感配置,后者用于声明持久存储需求。不同集群配置会用到不同辅助资源,列在 Owns 中不表示每套 RayCluster 都会创建它们。
Owns 关注下级资源。Pod 事件到来后,处理器沿 Pod 的 owner reference 查找 RayCluster,然后把这个 RayCluster 的标识放进队列。于是,Pod 的变化可以触发集群层面的检查。
这个行为能在 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,要等读到对象后再判断。
Pod 先触发 RayCluster 协调。RayCluster Controller 重新观察 Pod,计算集群状态并写回 API。RayJob Controller 通过自己注册的 Owns(RayCluster) 观察到变化,再检查 demo-job 是否可以进入提交阶段。这是两次不同 Controller 的处理,中间由资源状态衔接。
本例中,calculateStatus 还要确认资源维护没有报错、Pod 总数符合预期,且 CheckAllPodsRunning 检查通过。这个函数不仅检查 Pod 的 phase 是否为 Running,也会检查已有的 Ready 条件;只要某个 Ready 条件不是 True,就返回 false。因此,“Pod 已 Running”不能单独作为集群就绪的依据。这里特意说“已有的条件”:此函数没有额外拒绝缺少 Ready 条件的 Pod,不能把实现改述成“必须找到 Ready=True 才通过”。
上面那组 predicate 还有一个容易读错的地方:它写在 For(RayCluster) 上,只筛选这一条监听路径。对 RayCluster 的更新,只改 status 通常不会改变 generation,因此不会仅凭 generation 条件触发它自身的协调;label 和 annotation 的变化则有另外两个条件。Owns(Pod) 没有套上同一组过滤条件,Pod 状态变化仍然可以触发集群协调。RayJob 对下级 RayCluster 的监听也独立配置。
判断事件为什么没有入队时,应当沿当前 Controller 的这条监听路径检查过滤条件。另一个 Controller 或另一条监听上的过滤器,可能与它无关。
3. 进入 Reconcile 后,重新读取对象
从队列取出请求时,触发它的 Pod 可能已经又更新过几次,甚至已经删除。协调逻辑必须依据当前能读到的对象作判断。
RayCluster 的 Reconcile 开头就进行了读取。下面是连续节选:
1 | |
读到 NotFound,后面的分支会结束本轮处理;其他读取失败按错误处理。这也处理了对象已删除、旧请求却还留在队列里的情况。
进入 rayClusterReconcile 后,代码会检查资源是否交给外部 Controller 管理,验证配置,处理删除流程,再按当前状态维护集群资源。不同配置会走不同分支。本例没有删除请求,也不启用自动扩缩容,重点是 head Service、head Pod、worker Pod 和集群状态。
假设配置要求两个 worker,本轮只观察到一个,Controller 会结合已有 Pod 和尚未被缓存观察到的创建记录,判断是否需要补建。如果 Pod 数量已满足,但尚未全部达到就绪条件,就继续观察并更新 status。发出 Pod 创建请求以后,它不必一直占着这个协调协程等待容器启动,后面的事件或延时请求会带来下一轮检查。
spec 与 status 在这里分别承担不同作用。spec 保存使用者的期望,status 保存 Controller 已观察到的集群情况。实际 Pod、Service 等对象也参与判断。只读某一个 status 字段,不能代替完整的资源检查。
由于每轮都会重新比较期望和当前状态,多次相似事件可以合并,不必逐条重放历史通知。不过,同一个对象仍会反复检查,写入动作也得考虑重复调用。幂等要求重复执行同一操作不会额外产生一次效果;控制循环本身不会替被调用的外部接口实现这个要求。
4. 返回以后,谁安排下一轮
除了等待资源事件,Controller 还可以主动要求稍后重查。
RayJob 等待集群就绪时,或者轮询 Ray Jobs API 时,相关变化不一定立刻表现成一个它所监听的 Kubernetes 事件。因此,代码会返回 RequeueAfter,请求稍后再处理同一个对象。
RayJob.Reconcile 的收尾部分包含下面的代码。两段之间省略了一次指标更新:
1 | |
这两处都返回了相同的延时,实际是否采用,还要看 error。框架先处理错误:普通 error 非空时走限速重试,并忽略 RequeueAfter。这里由 rate limiter 决定重试的等待时间,避免失败请求立即反复执行。
在 controller-runtime reconcileHandler 中,成功且设置延时的分支使用下面两行连续代码:
1 | |
Forget 先清除这次请求的限速重试记录,下一行再安排延时入队。这里使用的是本版本优先级队列的 AddWithOpts;旧版本里常见的 AddAfter 并不是这份源码的调用。
| 返回结果 | 框架如何处理 |
|---|---|
| 普通 error 非空 | 记录错误,按 rate limiter 安排重试;忽略返回值中的延时 |
TerminalError |
本次错误不触发自动限速重试;后续资源事件仍可能再次触发处理 |
error 为空,RequeueAfter > 0 |
清除限速记录,安排延时请求 |
error 为空,RequeueAfter <= 0 且 Requeue: true |
走限速重排路径;这是该版本仍兼容的旧字段 |
| error 为空,空 Result | 本轮不主动安排重查;后续事件仍可以再次入队 |
RequeueAfter 指定的时间到了,请求还要等可用的协调协程;期间如果来了新事件,同一个对象又可能更早被处理。因此,“每隔两秒准确执行一次”并不是这里的调度语义。Reconcile 每次处理一轮,然后返回,不是每个对象长期占用一个循环。
5. 资源已经变化,却没有继续处理
假设最后一个 worker 已经就绪,作业却迟迟没有提交,可以沿前面的链路倒查。
首先确认 Pod 位于 Operator 的监听范围内,并且相关标签没有被改坏。然后看 owner reference:如果 Pod 没有指向正确的 RayCluster,Owns(Pod) 就无法把事件映射到预期对象。再检查 RayCluster Controller 的日志,确认它是否进入了这一轮,以及是否在校验、删除或资源维护分支提前返回。
如果 RayCluster 的状态已经更新,下一步看 RayJob 与集群之间的关系,再看 RayJob 当前状态是否满足提交条件。集群 Ready 以后仍可能在获取 Dashboard 地址或创建 submitter Job 时遇到错误,这些问题不会因为 Pod 已经就绪而自动消失。
| 检查位置 | 要确认的事情 |
|---|---|
| watch 与缓存 | namespace、标签筛选、RBAC、初始同步是否正常 |
| handler 与 predicate | 事件是否通过当前路径的过滤,能否找到正确的 owner |
| 队列与协调协程 | 请求是否还在等待,是否存在大量重试或慢调用 |
| Reconcile 分支 | 当前声明和状态让代码走到了哪里,是否有提前返回 |
| 返回结果 | 是等待下一次事件、安排延时,还是返回错误后限速重试 |
这些检查围绕同一条路径:Pod 更新进入缓存和事件处理器,RayCluster 请求入队,协调结果写回集群状态,RayJob 再接着处理。下一篇把这条路径放到 100 个 RayCluster 同时变化的场景里,分析哪些检查可以并行,以及同一个对象为什么不会被本地队列同时交给两个协调协程。
参考资料
- KubeRay 源码基线(6bf05eb17a3e)
- controller-runtime 源码基线(v0.24.1)
- KubeRay:ray-operator/main.go
- KubeRay:ray-operator/internal/managercache/cache.go
- controller-runtime:pkg/internal/source/kind.go
- controller-runtime:pkg/client/client.go
- KubeRay:ray-operator/controllers/ray/raycluster_controller.go
- controller-runtime:pkg/builder/controller.go
- controller-runtime:pkg/handler/enqueue_owner.go
- KubeRay:ray-operator/controllers/ray/rayjob_controller.go
- controller-runtime:pkg/internal/controller/controller.go
- KubeRay:ray-operator/controllers/ray/utils/util.go
- Kubernetes RBAC