KubeRay 源码阅读:一个 Operator 如何管理 100 个 RayCluster
第二篇追到了请求入队和一轮协调的返回。现在把第一篇的 demo-job 放大来看:每个作业创建一套专用 RayCluster,包含一个 head 和 compute 组的两个 worker。如果同一个 Kubernetes 集群里同时有 100 套这样的 RayCluster,Operator 怎样安排管理工作?每套集群各占一个协程,还是大家排队处理?配置中的协调并发数设为 2,又限制了什么?
要回答这个问题,先把协调与计算分开。协调并发数限制一个 Controller 同时执行多少轮 Reconcile;函数返回后,Ray 集群仍继续运行,协调协程则可以去检查另一个对象。至于 100 个集群能否得到及时管理,还要看事件产生得多快、每轮处理多久,以及 Kubernetes API 能承受多少请求。
本文基于 KubeRay 6bf05eb17a3e,依赖版本按该提交的 go.mod 核对为 controller-runtime v0.24.1、client-go v0.37.0。以下时间线用于解释源码行为,100 个集群是讨论场景,没有进行容量实测。
1. 两个协调协程怎样照看 100 个集群
每个 Controller 有自己的工作队列。队列里的工作项用来标识待协调对象,标准 reconcile.Request 只包含 namespace 和 name,不携带完整对象或原始事件。下文的“键”就是队列用来区分对象的这个标识。沿用前两篇的名称,default/demo-job-abcde 表示下一轮应读取并检查哪一个 RayCluster。
KubeRay 的 --reconcile-concurrency 默认值是 1。RayCluster Controller 把它传给 MaxConcurrentReconciles,controller-runtime 据此启动协调协程。下面节选自其 Controller 启动代码,仅省略循环内的两行解释性注释:
1 | |
这段循环里,每个协程取一个工作项,完成一轮协调,再取下一个,并不固定服务于某个 RayCluster。代码里的 worker 是 Controller 的协调协程;后文提到执行用户程序的 Ray worker 时,要注意区分。
把并发数设为 2,可以画出下面的调度过程。
为避免在时间线上反复写长名称,图中用 A、B、C、D 代指不同 RayCluster 的队列键,其中 A 可以取前文的 default/demo-job-abcde。时间向右,条块表示一轮 Reconcile,长度仅作示意。A、B 的两轮处理可以重叠;A 返回后,空出的协程可以接手 C。下方 Ray 集群的运行时间可以跨越多轮协调,不受这些条块的起止约束。
等待 Pod 就绪时,Controller 可以返回 RequeueAfter,让后续轮次继续检查。返回后协调协程就释放了;如果本轮还在等待 HTTP 响应,它仍占着一个并发名额。因此,提高协调并发数可能缩短其他对象的排队时间,却不会直接提高 Ray task、actor 的执行并发,也不能据此推导可承载的集群总数。
稳定集群也有管理开销。即使没有新的资源事件,队列仍可能出现待处理项。RayCluster 主流程末尾会安排周期性重新检查,间隔可由环境变量配置。这些检查通常不需要每轮都创建资源。
2. 同一个对象为什么不会被两个协程同时处理
假设协程 1 正在协调 A,它的两个 worker Pod 又接连发生变化,事件都映射成 A 的工作项。此时协程 2 即使空闲,也不能从同一队列再取到 A,启动重叠的一轮。
这部分要先确认实际用了哪种队列。controller-runtime v0.24.1 的 NewTypedUnmanaged 默认创建 priorityqueue;这份 KubeRay 配置没有覆盖 NewQueue 或关闭 UsePriorityQueue。client-go 是 Kubernetes 的 Go 客户端库,它的 v0.37.0 版本中,普通 workqueue 仍有 dirty、processing 集合,但不能拿它们当作这里实际使用的数据结构。
priorityqueue 用按键索引的 items 合并重复入队,并分别管理已到期的 ready 项和延迟等待的 waiting 项。这里的 ready 表示队列项已经可以等待分配协调协程,与 Pod 或 RayCluster 是否就绪无关。交出工作项时,它会跳过已锁定的键,然后锁定本次选中的键。以下是 handleReadyItems 的连续节选:
1 | |
队列内部按顺序遍历候选项,这里的 return true 表示继续寻找下一个可交出的项。Controller 处理函数用 defer c.Queue.Done(obj) 在本轮结束时通知队列,Done 再解除这个键的锁定。
图中 A 正在处理时,新到的 A 可以登记为待处理项,但要等前一轮 Done(A) 后才可再次交出。反复登记同一个待处理键会合并;优先级取较高值,等待期限取较早值。队列按优先级选择已到期的项,同优先级再按入队次序排序,不能把整张图读成严格的先进先出(FIFO)。
重新排队还会调整次序:延迟项到期进入 ready 时重新取得次序,提升优先级的项排到新优先级已有项之后。如果 A 已安排延迟重查,又因新事件立即入队,合并结果可以提前变为可处理。RequeueAfter 因而不保证“这段时间内绝不会再协调”,到期也不保证立刻拿到协调协程。
锁定只发生在同一 Controller 实例的同一队列里。其他 Controller、其他进程和用户仍可以修改这个 Kubernetes 对象。由于重复事件会合并,协调函数还必须重新读取状态,不能依赖每条 Pod 事件都触发一次独立调用。
如果积压主要来自 A 一个对象,增加协程也无法把 A 的两轮协调拆开并行。其他协程可以处理 B、C,A 仍要等待当前轮返回。遇到这种情况,应检查 A 为什么反复触发或单轮耗时很长,不能仅凭事件很多就把并发参数调高。
3. 已经串行,为什么还需要 expectations
同一个对象串行处理后,是否就能避免重复创建 Pod?继续看 demo-job 的 compute 组:期望两个 worker,当前缓存只看到一个。第一轮补建 Pod,API 返回成功;第二轮紧接着开始,informer 缓存里却可能仍只有原来的一个。
两轮函数完全串行,仍可能根据旧列表重复补 Pod。原因在读写路径:Manager 客户端的普通 Get、List 可以从缓存读取,创建、更新则请求 API Server。写入成功并不代表对应 watch 更新已经被本地缓存消费。启动时的缓存同步和 handler 初始事件同步,也不会让之后的每次读都与 API Server 同步。
KubeRay 用 scale expectations 记录尚待观察到的 Pod 创建、删除结果。在 createWorkerPod 中,登记发生在创建返回成功之后。下面只省略失败分支中的事件记录:
1 | |
下一轮处理某个 worker group 之前,会调用 IsSatisfied。如果该组还有未满足的预期,本轮跳过这一组;其他组可以继续检查。head 也有对应检查,使用空字符串作为组标识。
图中两个时刻之间,API 已经创建 Pod,缓存却还没看到结果。预期记录让下一轮知道还有创建结果尚待观察。它按 namespace、RayCluster 和 group 查询,所以 A 的 Pod 暂时不可见,不会阻止 B 继续协调。
判断条件比“等待创建事件”更具体。isPodScaled 按名字调用客户端的 Get 读取 Pod;这个客户端来自 Manager,普通 Pod 读取仍可能经过缓存,并没有在这里强制直读 API Server。下面保留完整的两个分支,只省略创建分支中说明极端情况的注释:
1 | |
创建预期在能读到 Pod 时满足,不要求它已经 Ready;读不到时,超过源码设置的 30 秒也会放行。删除预期要求读取返回 NotFound,没有同样的超时。Pod 已有删除时间戳、但仍可被查询到时,这条预期仍未满足;命令行此时常显示 Terminating,它不是 Pod 的一种 phase。
这些记录通过 client-go 的线程安全 Indexer 保存在进程内。Indexer 是支持按索引查找的内存存储,它的锁保护内部访问;expectations 则让控制器在缓存滞后时少做重复的扩缩容判断。记录会随进程退出而丢失,创建预期还有超时,所以它不能保证跨重启“只执行一次”。如果创建实际成功、响应却丢失了,前面的登记甚至可能根本没有执行。
4. 多个 Controller、namespace 和 Kubernetes 集群各是什么边界
demo-job 创建 RayCluster 后,RayJob Controller 可以查询作业状态,RayCluster Controller 同时检查 worker Pod。二者在 main.go 中分别注册,拥有各自的队列和协调协程,共享 Manager 提供的基础设施。
同一个 config.ReconcileConcurrency 分别传给 RayCluster、RayJob、RayService。因此设为 2 时,这三个 Controller 各自最多有两轮已经开始、尚未返回的协调,合计可以达到六轮;后文将这种状态称为“在途”。这只是三个核心 Controller 的理论上限,实际还受工作项是否可处理等条件限制。它不是整个 Operator 共用的两个名额。
这些协调共享进程资源,也可能读写同一个下级 Kubernetes 对象。不同键虽然可以同时处理,具体操作仍可能冲突。排查时要继续看它们改的是哪个对象、哪个字段;若访问了共享内存,再检查相应的同步方式。
图中的队列位于 Controller 框内;namespace 框表示资源所在范围。这里另设一个同名对象的对比例子:team-a/demo-job-abcde 与 team-b/demo-job-abcde 分属两个命名空间,是不同键,可以由同一个 RayCluster Controller 并发处理。主案例仍在 default 中,并没有迁移到这两个命名空间。namespace 不会自动获得独立队列、专属协程或公平份额。
KubeRay 的 --watch-namespace 可限定一个或逗号分隔的多个 namespace,空值表示监听所有 namespace。这个配置约束缓存观察范围;实际访问权限还取决于 RBAC。
100 个 RayCluster 可以属于同一个 Kubernetes 集群。这里的启动入口只用一份 rest.Config 创建一个 Manager;这份客户端配置指定 Kubernetes API 地址、访问身份等信息。增加 namespace 不会让它连接第二个 Kubernetes API。管理另一个 Kubernetes 集群,需要为它配置对应的 Operator 管理实例。
增加 Operator 副本还涉及选主。leader election 从一组实例中选出负责协调的主实例,其余等待接管。此基线默认开启这项机制,所以共享同一选举锁的副本不能简单视为并发数相加。若多个独立实例的观察范围重叠,各自队列的键锁也无法互相排斥。选举和接管留到第五篇展开。
5. 对象被别人改了,版本检查怎样起作用
即使同一个 RayCluster 在本队列中串行处理,用户仍可能同时修改它的 spec。Controller 读到版本 A 后准备写 status,API Server 中的对象已经变为版本 B,这时就需要乐观并发检查:读取时不排斥别人修改,写入时再确认对象是否仍是读到的版本。
resourceVersion 是对象的版本标记。客户端应原样携带它,不自行推算或递增。普通带版本的 Update 提交旧版本时,API Server 的版本比较会拒绝更新并返回 Conflict。RayCluster 的 updateRayClusterStatus 使用 r.Status().Update(ctx, newInstance),把错误交回主流程。
主流程在选择错误后执行以下判断。这是 Reconcile 返回处的连续节选:
1 | |
同时返回错误和等待时间时,要继续检查框架。reconcileHandler 优先处理非空错误:普通错误通过 AddWithOpts 的 RateLimited 选项重新入队,忽略 RequeueAfter;终止型错误有单独分支。只有没有错误时,正的 RequeueAfter 才用于延迟入队。
因此,版本冲突后会通过后续协调重新读取和计算,不能拿原来的旧对象原样重试并期待成功。缓存仍可能暂时落后,重试也可能再次冲突。版本检查保护一次对象更新,不会把 Pod 创建与 RayCluster status 更新组成事务。API 还支持只描述部分修改的 Patch 请求,其冲突条件取决于具体写法,不能一概套用这里的 Update 规则。
比如本轮已经创建 Pod,随后写 status 发生冲突,先前成功的 Pod 创建不会被回滚。下一轮应结合重新读到的声明、Pod 和预期记录继续判断。排障时只看最后一条 Conflict 日志,容易漏掉本轮前半段已经完成的资源操作。
6. 怎样判断该不该提高并发数
同样是 100 个集群,稳定运行与同时扩容给 Operator 带来的负载差别很大。决定是否提高并发数之前,需要按 Controller 看请求在哪里等、处理多久,以及错误是否增多。
下表列出 controller-runtime 提供的 Controller 和队列指标。
| 指标 | 要判断的问题 |
|---|---|
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 的集群数。
如果协调协程长期占满、ready 项持续积压,而单轮耗时和 API 请求表现稳定,提高并发值得验证。若 API 已经变慢或限流等待增加,继续加协程可能增加争用。KubeRay 还注册了 rest_client_request_duration_seconds 和 rest_client_rate_limiter_duration_seconds,分别帮助观察请求耗时和客户端限流等待;请求 QPS、Burst 在 Manager 创建前配置,分别约束客户端的持续请求速率与突发额度。
验证容量时,应固定每个 RayCluster 的 Pod 规模和变更负载,比较调整前后的排队耗时、错误速率及集群就绪耗时。只有“管理了 100 个 CR”这个数字,还无法判断它们是否得到了及时处理。另一个问题是,一轮协调如果在创建资源后中断,状态还没来得及写回,下一轮根据什么接着做?下一篇沿这个窗口分析恢复过程。
若 Pod 已经创建,却因节点资源不足而无法启动,需要继续分析 Pod 调度与作业间资源竞争。第六至第九篇会引入 Volcano,区分这类等待与本篇的 Controller 工作队列积压。
参考资料
- KubeRay 源码基线(6bf05eb17a3e)
- KubeRay:ray-operator/go.mod
- controller-runtime:pkg/reconcile/reconcile.go
- KubeRay:ray-operator/main.go
- KubeRay:ray-operator/apis/config/v1alpha1/defaults.go
- KubeRay:ray-operator/controllers/ray/raycluster_controller.go
- controller-runtime:pkg/internal/controller/controller.go
- controller-runtime:pkg/controller/controller.go
- client-go:util/workqueue/queue.go
- controller-runtime:pkg/controller/priorityqueue/priorityqueue.go
- controller-runtime:pkg/client/client.go
- controller-runtime:pkg/internal/source/kind.go
- KubeRay:ray-operator/controllers/ray/expectations/scale_expectations.go
- client-go:tools/cache/thread_safe_store.go
- apiserver:pkg/registry/generic/registry/store.go
- controller-runtime:pkg/internal/controller/metrics/metrics.go
- controller-runtime:pkg/internal/metrics/workqueue.go
- controller-runtime:pkg/controller/priorityqueue/metrics.go
- KubeRay:ray-operator/controllers/ray/metrics/client_go_metrics.go
- Kubernetes:Pod 生命周期与 Terminating 显示