第 3 篇|KubeRay 并发:一个 Operator 如何管理 100 个 RayCluster

本文最后更新于:2 天前

前言

第二篇说明了一个 RayCluster 怎样得到一次检查机会,以及本轮结束后怎样继续处理。现在把范围扩大到同一个 Kubernetes 集群里的 100 个 RayCluster:Operator 怎样在它们之间分配处理机会?配置中的协调并发数设为 2,究竟限制了什么?

下面先看不同 RayCluster 怎样共享协调协程,以及同一个 RayCluster 的多次请求为什么不会并行执行;再进入队列内部,说明请求如何合并、等待和分配。取得请求之后,Controller 还要面对缓存同步延迟和其他写入者,这两类问题决定了资源操作怎样安全接续。

多个集群怎样共享协调协程

先看 RayCluster Controller。它有自己的工作队列,队列中的请求标识需要检查哪套集群。例如 default/demo-job-abcde,斜杠前是命名空间,后面是对象名;源码中的标准 reconcile.Request 保存这两个字段。队列用这个标识区分对象,后文把它称为“键”。

负责从队列取请求并调用 Reconcile 的执行单元叫协调协程。一轮协调会读取资源、判断差异,完成当前能做的创建、更新等操作,然后返回。协程可以接着处理另一套集群,并不固定属于某个 RayCluster。

KubeRay 的 --reconcile-concurrency 默认值是 1。RayCluster Controller 把它传给框架的 MaxConcurrentReconciles,决定启动多少个这样的协程。下面是 controller-runtime 启动协调协程的代码,仅省略循环内的两行解释性注释:

1
2
3
4
5
6
7
8
wg.Add(c.MaxConcurrentReconciles)
for i := 0; i < c.MaxConcurrentReconciles; i++ {
go func() {
defer wg.Done()
for c.processNextWorkItem(ctx) {
}
}()
}

循环中的 go func() 启动一个 Go 协程;每个协程反复执行 processNextWorkItem,取出请求并完成一轮处理。外面的循环次数就是配置的并发数。wg 用于等待这些协程结束,不参与请求去重。源码里这类执行单元也称 worker,它负责管理资源;Ray worker 则在 Ray 集群里执行用户计算。

假设并发数设为 2,集群 A、B、C、D 都需要检查。其中 A 仍是 default/demo-job-abcde,其余字母代表另外三套 RayCluster。开始时两个协程可以分别处理 A、B;A 的这一轮先结束,空出的协程就可以接手 C。

图 1:两个协调协程轮流协调多个 RayCluster
图 1:两个协调协程轮流协调多个 RayCluster 查看原图

图中时间向右,条块表示一轮 Reconcile,长度仅作示意。A 的检查结束后,它的 Ray 程序仍然可以继续运行。因而,两个协调协程可以轮流管理很多集群,配置值 2 不代表只能运行两个 RayCluster,也不限制 Ray task、actor 的计算并发。

是否占用这个名额,要看当前函数有没有返回。如果 Controller 正在等待一次 HTTP 请求的响应,这一轮尚未结束,仍占一个名额。如果已经发现 Pod 尚未就绪,返回 RequeueAfter 要求稍后再查,本轮就结束了,延时等待由队列安排,协程可以去处理其他集群。

同一个集群的请求怎样接续

A 和 B 可以同时处理,同一个 A 的两轮协调则按顺序进行。这样,针对 A 的一次数量判断和资源操作结束后,下一轮才会开始,避免两轮根据各自读到的 Pod 列表同时补建 worker。

这里要分别记录两件事:A 是否还有一次检查等待执行,以及 A 是否已经有一轮正在执行。下面暂时忽略延时和优先级,只看立即请求经过队列整理后的变化。

时刻 发生的事情 A 的待处理记录 A 是否正在处理
第一次变化 登记一次 A 的检查请求 有一份 否
开始协调 协程取走这份请求 无 是
处理中再次变化 为 A 登记下一次检查 有一份 是
又来一次变化 合并到已有待处理记录 仍是一份 是
本轮结束 解除 A 的执行锁定 保留,等待分配 否
再次分配 协程取走下一次请求 无 是

表中间两行描述的是:A 正在处理,同时还有一份 A 等待处理。当前轮可能已经读过 Pod,后面发生的新变化就留给下一轮检查。这份请求会先保留在队列中,等当前轮结束后再交给协调协程。

图 2:同一个键重复入队、锁定与再次处理
图 2:同一个键重复入队、锁定与再次处理 查看原图

等待 A 的时候,空闲协程仍可以接手 B、C 等其他对象。同一个待处理请求反复到来时会合并;每一轮实际开始后,再读取当时能看到的资源状态。队列不需要为每条 Pod 通知安排一轮独立调用。

如果 A 处理期间没有留下新的请求,本轮结束后也就没有下一份 A 可以分配。后来再次发生变化,仍可以重新登记。去重的范围是当前的待处理记录。

这也说明了提高并发数的作用范围:如果积压来自许多不同集群,更多协程可以让它们同时得到处理;如果反复变化的主要是 A,A 仍要一轮接一轮执行,其他协程不能把它的两轮协调拆开并行。

队列怎样实现这些安排

前面的表说明了请求在等待和执行之间怎样变化。下面把负责这些动作的协程分开,再看它们如何共享队列记录、交接任务。

Controller 使用 controller-runtime 的优先级队列 priorityqueue。创建入口 NewTypedUnmanaged 默认选择它;KubeRay 沿用这一配置,没有通过 NewQueue 替换队列,也没有关闭 UsePriorityQueue。

谁安排任务,谁执行协调

继续采用协调并发数为 2 的例子。一个 RayCluster Controller 的请求处理路径涉及以下协程:

所在位置 分工与数量 具体工作
队列内部 1 个整理请求的协程 处理输入缓冲,合并重复请求,登记延时和优先级
队列内部 1 个管理延时的协程 等待延时到期,把对应项转入可处理集合
队列内部 1 个分配任务的协程 选择可执行对象,交给正在取任务的协调协程
Controller 2 个协调协程 向队列取任务,取得对象后调用 Reconcile

队列创建时分别启动 handleAddBuffer、handleWaitingItems 和 handleReadyItems,对应前三行。最后一行由 Controller 按 MaxConcurrentReconciles 启动。这个例子是三个队列后台协程加两个协调协程,配置中的并发数只控制后者;日志、指标等还有其他辅助协程,所以表格不是进程的协程总数。

前三类协程负责维护请求和分配条件;取得任务后的协调协程负责读取集群资源、创建 Pod、写回状态。它们都在同一进程中工作,队列内部通过锁保护共享记录,通过通知唤醒需要继续工作的协程。

请求怎样进入输入缓冲

提交请求时,调用方通过 AddWithOpts 交出对象标识和处理选项;选项可以表达延时、优先级或限速重试。入口先把它们存入输入缓冲,再通知后台整理。下面是接收请求的连续节选:

1
2
3
4
5
6
7
8
w.addBufferLock.Lock()
w.addBuffer = append(w.addBuffer, bufferItem[T]{
opts: o,
items: items,
})
w.addBufferLock.Unlock()

w.notifyItemAddedToAddBuffer()

addBuffer 是刚收到的请求暂存区,opts 保存选项,items 保存这次提交的对象标识。追加前后加锁和解锁,是为了让多个提交方安全地访问这个缓冲。最后的 notifyItemAddedToAddBuffer 发出“有新请求”的通知。

此处先追加请求,按对象标识合并的工作留到整理阶段。例如 A 连续提交两次,缓冲里可以同时留着这两次请求。调用 AddWithOpts 返回时,请求已经被接收,随后再参与合并。

输入缓冲暂存的是“检查 A”的请求。后文还会用到对象缓存:它由 informer 维护,保存 RayCluster、Pod 等资源的本地记录,供 Reconcile 读取。两者分别服务于任务安排和资源观察。

整理协程:合并请求并通知后续处理

整理协程负责把输入缓冲中的请求变成待处理记录。它等待 itemAddedToAddBuffer 通知,运行下面的 handleAddBuffer 循环:

1
2
3
4
5
6
7
8
9
10
11
12
13
func (w *priorityqueue[T]) handleAddBuffer() {
for {
select {
case <-w.done:
return
case <-w.itemAddedToAddBuffer:
}

w.lock.Lock()
w.lockedFlushAddBuffer()
w.lock.Unlock()
}
}

可以按一次通知读这段循环:select 先等待;收到 itemAddedToAddBuffer 通知后,取得队列锁,调用 lockedFlushAddBuffer 整理缓冲,再解锁并回到等待。若收到队列关闭信号 done,这个协程就退出。

整理时,lockedFlushAddBuffer 取走当时已有的一批缓冲记录,交给 lockedAddWithOpts 加入或更新待处理项。一次通知可以带来一批处理,没有固定刷新周期。输入缓冲将接收请求和维护队列结构分开,提交方不必一直持有队列主锁做后续整理。

整理后的待处理项存放在按对象标识索引的 items 中。这里的 items 是整个队列的待处理索引,与前面缓冲记录中“本次提交了哪些对象”的同名字段处在不同层次。

以 A 为例:索引里没有 A,就建立一项;已经有 A,就把新要求合并到现有项。前面表中“仍是一份”对应的就是这一步。索引只保存尚待处理的工作,正在执行的那一轮另有记录。

请求有时需要立即处理,有时要求稍后再查。队列另外维护两个集合,按时间条件组织这些待处理项:

集合 保存什么 能否交给协调协程
waiting 尚未到期的延时项 先等待时间条件满足
ready 无须延时或已经到期的项 还要看对象是否正在处理、是否有空闲协程

这里的 ready 表示队列项的时间条件满足,与 Pod 或 RayCluster 的 Ready 状态没有关系。items、waiting 和 ready 是同一批待处理工作的不同组织方式,不是依次执行的三份独立任务。

整理完成后,还要通知谁可以继续工作。下面是 lockedAddWithOpts 收尾的连续节选;两个布尔变量记录这一批处理是否新增 ready 项、是否新增或更新 waiting 项:

1
2
3
4
5
6
if readyItemAdded {
w.notifyReadyItemOrWaiterAdded()
}
if waitingItemAddedOrUpdated {
w.notifyWaitingItemAddedOrUpdated()
}

新增 ready 项时通知分配协程;新增或更新 waiting 项时通知延时协程。一批请求可能同时产生这两类通知。通知只表示“共享记录发生变化,请重新检查”,具体对象仍保存在队列中。整理协程随后释放队列锁,回到循环等待下一次输入通知,不会在这里执行 Reconcile。

延时协程:等待到期并转入可处理集合

假设 A 要求 10 秒后再查,整理阶段会把它放进 waiting。延时协程 handleWaitingItems 负责等待这个处理时间;即使所有协调协程都在忙,它也可以独立更新队列中的到期状态。

它等待三种信号:队列关闭、新增或更新了延时项,以及先前安排的到期信号。下面保留完整函数。ReadyAt 是项的到期时间,toMove 暂存本轮发现的已到期项,nextReady 是下一次时间信号;Ascend 按 waiting 的顺序遍历,最早到期的项在前面。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
func (w *priorityqueue[T]) handleWaitingItems() {
blockForever := make(chan time.Time)
var nextReady <-chan time.Time
nextReady = blockForever

for {
select {
case <-w.done:
return
case <-w.waitingItemAddedOrUpdated:
case <-nextReady:
nextReady = blockForever
}

func() {
w.lock.Lock()
defer w.lock.Unlock()

var toMove []*item[T]
w.waiting.Ascend(func(item *item[T]) bool {
readyIn := item.ReadyAt.Sub(w.now()) // Store this to prevent TOCTOU issues
if readyIn <= 0 {
toMove = append(toMove, item)
return true
}

nextReady = w.tick(readyIn)
return false
})

// Don't manipulate the tree from within Ascend
for _, toMove := range toMove {
w.waiting.Delete(toMove)
toMove.ReadyAt = nil

// Bump added counter so items get sorted by when
// they became ready, not when they were added.
toMove.AddedCounter = w.addedCounter
w.addedCounter++

w.metrics.add(toMove.Key, toMove.Priority)
w.ready.ReplaceOrInsert(toMove)
}

if len(toMove) > 0 {
w.notifyReadyItemOrWaiterAdded()
}
}()
}
}

可以沿一次唤醒看完整流程:

  • 初始的 blockForever 不会发出时间信号,协程先等待延时项的新增或更新通知;队列关闭则直接退出。
  • 唤醒后取得队列锁,从最早到期的项开始检查。已到期的加入 toMove;遇到第一个尚未到期的项,就按它剩余的等待时间设置 nextReady,停止本轮遍历。若新加入的任务更早到期,这次检查也会调整下一次等待。
  • 遍历结束后,将 toMove 中的项从 waiting 删除,清除 ReadyAt,更新排序次序和指标,再加入 ready。代码先收集再移动,把遍历和修改树分成两步。
  • 只要移动了项,就通知分配协程,随后释放锁并回到外层循环。进入 ready 后,该项可以参与任务选择;具体何时交出,还要看空闲协调协程和对象的执行锁定状态。

例如,A 到期后会从 waiting 转入 ready;如果 A 的上一轮还没结束,它仍不能立即执行。没有延时的 B 则在整理阶段直接进入 ready,整个过程不需要经过延时协程。

协调协程主动取任务,队列响应这个请求

整理请求和交出任务由不同的后台协程负责。要理解交接,先看接收任务的一端:协调协程完成上一轮后,会再次调用队列的取任务方法。Controller 的 processNextWorkItem 调用的是 GetWithPriority,同时取得对象标识和优先级;队列的普通 Get 也会转调这个方法。

如果暂时没有任务,协调协程就在取任务方法中等待。GetWithPriority 先登记一个等待者,再通知任务分配逻辑;下面是连续节选:

1
2
3
4
5
w.lock.Lock()
w.waiters++
w.lock.Unlock()

w.notifyReadyItemOrWaiterAdded()

waiters 记录有多少协调协程已经来取任务、尚未拿到结果。后面的 select 等待队列关闭,或者从 w.get 通道接收一个任务。这里的等待不会不断轮询,也不会执行 Reconcile。

例如,协调协程 1 调用取任务方法后停在通道接收处;A 的请求整理好以后,分配协程确认有人等待,便可以选择 A,通过通道发送。协调协程 1 收到 A,取任务方法返回,才开始执行 Reconcile(A)。这是同一次交接的两端:协调协程主动申请任务,队列内部通过发送响应它,不会因此额外启动一个协调协程。

分配协程:选择对象并交给协调协程

分配协程 handleReadyItems 把 ready 中的对象与已经来取任务的协调协程配对。它的外层循环先等待通知,下面是连续节选:

1
2
3
4
5
select {
case <-w.done:
return
case <-w.readyItemOrWaiterAdded:
}

整理协程新增 ready 项、延时协程移入到期项、协调协程来取任务,以及 Done 解除对象锁定,都会通过 readyItemOrWaiterAdded 唤醒它。这些变化都可能让一次新的分配成为可能。

被唤醒后,它先取得队列锁,整理一次输入缓冲,再检查 waiters。如果没有协调协程等待取任务,就结束这次分配尝试、释放锁,回到外层循环等待通知。这里结束的是循环内的一次处理,分配协程本身仍在运行。有人等待时,它才继续遍历 ready,按照优先级和次序选择对象。

选择对象时,需要检查“正在处理”的记录。源码用 locked 保存这些对象的键。下面是跳过已锁定对象、锁定本次选中对象的连续节选:

1
2
3
4
5
6
if w.locked.Has(item.Key) {
return true
}

w.metrics.get(item.Key, item.Priority)
w.locked.Insert(item.Key)

假设当前候选项是 A。若 A 已在 locked 中,return true 表示继续遍历下一项,不会交出 A,也不会把它的待处理请求丢掉。若 A 没有锁定,则记录这次出队的指标,并将 A 加入 locked。

紧接着是本次交接的连续代码:

1
2
3
4
w.waiters--
delete(w.items, item.Key)
toDelete = append(toDelete, item)
w.get <- *item

这几行减少等待者计数,从待处理索引 items 移除 A,并通过 w.get 通道发送。前面调用取任务方法、正在等待接收的某个协调协程收到 A 后返回,随后开始 Reconcile(A)。toDelete 暂存已经交出的项,这一批遍历结束后再统一从 ready 删除。

如果还有等待者,分配协程继续找下一项;没有等待者就停止遍历。若 ready 中没有可交出的对象,例如剩下的键都已锁定,也会结束本次尝试。最后清理 ready、释放锁,回到外层循环等待新通知。它可以一次交出多个不同对象,但不会在本轮处理完后不停扫描空队列。

此时,待处理的 A 已取走,正在执行的 A 仍记录在 locked 中。若新的 A 到来,便可以重新建立待处理项,同时继续受这把键锁约束。

协调协程通过 defer c.Queue.Done(obj) 保证本轮退出时通知队列。Done(A) 清除 A 的执行锁定并通知任务分配逻辑。之前留下的 A 只有在时间条件满足、协程可用时,才会再次被交出。这样,待处理索引负责合并请求,执行锁定负责避免同一个对象的两轮协调重叠。

回到最初的 A 请求:提交方把它追加到输入缓冲,整理后成为一份待处理记录;有延时就先进入 waiting,到期后进入 ready。协调协程发出取任务请求,分配协程才从 ready 中挑选未锁定的对象并交出;分配路径也会顺手整理尚未处理的输入缓冲。A 的这一轮执行结束,由 Done 解除锁定,期间留下的下一份 A 才有机会继续。

延时和优先级怎样影响顺序

前面分别说明了立即请求和延时请求的去向。两种请求如果指向同一个对象,又会怎样合并?若 A 已经安排延时重查,期间又收到一次立即检查的请求,队列合并时会采用较早的处理时间;优先级有差异时采用较高值。因此,新事件可以让延时项提前进入 ready,但无法绕过 A 当前轮的执行锁定。

在可以分配的 ready 项之间,队列先看优先级,同优先级再看次序。延时项到期进入 ready 时重新取得次序;一项被提升优先级后,排在新优先级已有项之后。整个队列不能视为严格的先进先出(FIFO)。

RequeueAfter 表达的是本轮对后续检查的安排。新事件可能使下一轮提前发生,到期后也可能继续等待空闲协程,因此它不保证精准的固定周期。

两轮协调之间,为什么还需要记录创建结果

队列已经让 A 的两轮协调串行执行,但下一轮能否看到上一轮的写入,还取决于资源读取路径。继续看 A 的 compute worker 组:这个组要求两个同类 worker Pod,当前本地缓存只看到一个。

API 已经创建,缓存还没有看到

第一轮向 API Server 请求补建一个 Pod,并收到了成功响应。此时 API Server 中已经有新 Pod,但 Operator 本地缓存可能还没收到对应的监听更新。若第二轮马上读取 Pod 列表,仍可能只看到原来的一个。

观察位置 创建成功、缓存尚未更新时看到的情况
API Server 已经保存新 Pod,实际已有两个
Operator 的本地对象缓存 仍只记录原来的一个 Pod
刚执行创建的协调逻辑 知道创建已经成功,可以把这个结果记下来

Manager 为 Controller 提供资源客户端。它的普通 Get、List 可以从对象缓存读取,创建和更新则请求 API Server。写入结果通过监听逐步反映到缓存,因此两条路径之间有同步时间差。启动时会完成缓存和初始事件同步,运行期间的更新则持续到达;即使两轮协调串行执行,后一轮读取时也可能还处在这段同步过程中。

记下已经发出、尚待观察的扩缩容结果

为了让下一轮知道“还有一次创建成功的结果尚未观察到”,KubeRay 保存一份扩缩容预期记录,源码称为 scale expectations。记录包含所属 RayCluster、worker 组、Pod 名称以及创建或删除动作。

登记位置在 createWorkerPod 中:创建请求返回成功之后,才把这个 Pod 记入预期。下面只省略失败分支中的事件记录:

1
2
3
4
5
if err := r.Create(ctx, &replica); err != nil {
// ... 省略 Warning 事件记录。
return err
}
r.rayClusterScaleExpectation.ExpectScalePod(replica.Namespace, instance.Name, worker.GroupName, replica.Name, expectations.Create)

代入本例,这份记录表达的是“已经为 A 的 compute 组成功创建了这个名字的 Pod,等待后续读取确认”。下一轮即使暂时仍只读到一个 Pod,也能先检查是否还有未观察到的创建结果。

图 3:Pod 写入成功与缓存观察之间的时间差
图 3:Pod 写入成功与缓存观察之间的时间差 查看原图

处理每个 worker 组之前,Controller 调用 IsSatisfied 检查该组的预期。下面是 reconcilePods 中 worker 组循环开头的判断:

1
2
3
4
if !r.rayClusterScaleExpectation.IsSatisfied(ctx, instance.Namespace, instance.Name, worker.GroupName) {
logger.Info("reconcilePods", "worker group", worker.GroupName, "Expectation", "NotSatisfiedGroupExpectations, reconcile the group later")
continue
}

continue 跳过本轮对这个组的 Pod 维护,循环继续检查其他组。协调协程不会停在这里等缓存,下一轮走到这个组时还会重新判断。预期按 namespace、RayCluster 和 group 查询,A 的创建结果暂时不可见,也不会阻止 B 的协调。head 有对应检查,用空字符串作为组标识。

每次检查预期,并清理已经通过的记录

IsSatisfied 先从预期存储中查出本组记录,再逐条调用 isPodScaled 判断。下面保留它的检查循环;items 是查出的记录列表,rp 是当前这条 Pod 创建或删除预期:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
for i := range items {
rp := items[i].(*rayPod)
isPodSatisfied := r.isPodScaled(ctx, rp)

if !isPodSatisfied {
return false
}

// delete satisfied item in cache
if err := r.itemsCache.Delete(items[i]); err != nil {
// Fatal error in KeyFunc.
panic(err)
}
}
return true

遇到未满足的记录就返回 false,已经满足的记录则从预期存储中删除。全部通过,或者已经没有记录时,返回 true,本组继续维护 Pod。清理预期记录表示这次操作已经完成确认,Kubernetes Pod 本身保持原状。

预期随着后续协调中的读取和判断逐步清理。具体条件在 isPodScaled 中:它按记录中的 Pod 名称调用客户端 Get,沿用 Manager 客户端的普通读取路径,因此仍可能从对象缓存获取结果。下面保留完整的两个分支,只省略创建分支中的注释:

1
2
3
4
5
6
7
8
9
10
11
12
13
switch rp.action {
case Create:
if err := r.Get(ctx, types.NamespacedName{Name: rp.name, Namespace: rp.namespace}, pod); err == nil {
return true
}
// ... 省略创建后迅速被删除的情况说明。
return rp.recordTimestamp.Add(ExpectationsTimeout).Before(time.Now())
case Delete:
if err := r.Get(ctx, types.NamespacedName{Name: rp.name, Namespace: rp.namespace}, pod); err != nil {
return errors.IsNotFound(err)
}
}
return false

创建分支在能读到 Pod 时返回 true,尚不要求它 Ready;读取不成功时,则比较登记时间加上 ExpectationsTimeout 是否早于当前时间。这个值在源码中设为 30 秒,超时也会返回 true,让 IsSatisfied 清掉该记录。

删除分支只有读取返回 NotFound 才返回 true。仍能读到 Pod,或者读取发生其他错误,都会返回 false。Pod 虽然已经标记删除,但只要仍可查询到,删除预期就没有满足。命令行中的 Terminating 表示删除中的显示状态,不是 Pod 的一种 phase。

预期在后续协调中怎样更新

创建预期保留 30 秒的等待窗口。源码注释给出的例子是:Pod 创建成功后,在下一轮检查前就被其他组件删除,Controller 因而没有机会读到它存在。后续检查发现等待已超时,就清理这条记录,恢复按当前 Pod 列表计算数量差异,决定是否补建。这里恢复的是本组的资源维护,Pod 是否就绪仍由状态检查判断。

等待时长在后续调用 IsSatisfied、检查到这条记录时计算,并没有单独的 30 秒定时器。新的 Pod 事件可以带来下一轮协调,RayCluster 主流程的正常收尾也会返回延时重查要求:ctrl.Result{RequeueAfter: time.Duration(requeueAfterSeconds) * time.Second}。这个分支默认使用 300 秒,可通过环境变量调整;错误或状态变化等路径可能安排更早的处理。因而,30 秒是预期检查使用的超时条件,实际检查时刻由下一轮协调决定。

删除预期以读到 NotFound 作为完成条件,常规检查中没有对应的超时放行。只要记录仍保留、读取尚未返回 NotFound,后续协调就继续暂缓这个组的 Pod 维护,同时处理其他组和其他集群。因此,这个组恢复扩缩容的时机取决于何时确认对象已经消失;等待较久时,也沿用同一条判断规则。

预期也会随所属资源的生命周期清理。例如,RayCluster 读取返回 NotFound 时会清空相关预期;满足条件的 Recreate 升级流程在删除 Pod 后也会清空记录。这些清理由各自的生命周期条件触发,常规删除预期仍按前述读取结果判断,等待时长本身不会触发强制放行或强制清理 Pod。

这些预期保存在 Operator 进程内,底层使用 client-go 提供的线程安全 Indexer,即支持索引查询的内存存储。它保护记录的并发访问,预期本身则帮助下一轮识别缓存尚未反映的操作结果。

expectations 用于减少同一进程前后两轮协调中的重复操作。进程退出后记录随之丢失;API 已创建成功但响应丢失时,也可能还没来得及登记。创建预期超时后,Controller 重新按读到的 Pod 列表判断数量,此时缓存仍可能处于同步过程中。重启后怎样识别已有资源,还要结合资源身份和状态管理,下一篇继续讨论。

并发限制在哪些范围内生效

到这里,队列控制处理顺序,扩缩容预期补充缓存尚未反映的操作结果。它们都属于 Controller 本地的管理机制。扩大到不同 Controller、不同进程和使用者同时操作时,还要区分限制的范围。

各个 Controller 的名额和管理范围

RayJob Controller 可以查询作业状态,同时由 RayCluster Controller 检查 worker Pod。二者在 main.go 中分别注册,各有队列和协调协程,共享 Manager 提供的客户端等基础设施。

同一个 config.ReconcileConcurrency 分别传给 RayCluster、RayJob、RayService。设为 2 时,每个 Controller 各有两个协调名额。如果三边都有足够的待处理请求,合计可以有六轮协调同时进行。

图 4:Controller、namespace 与 Kubernetes 集群的管理范围
图 4:Controller、namespace 与 Kubernetes 集群的管理范围 查看原图

图中的队列位于各自的 Controller 内。队列键包含 namespace,所以 team-a/demo-job-abcde 与 team-b/demo-job-abcde 是两个不同对象,可以并发协调。这只是同名对象的对比例子,主案例 A 仍在 default 命名空间。namespace 用来区分资源,不会自动分得独立队列、专属协程或公平份额。

KubeRay 的 --watch-namespace 可以限定一个或逗号分隔的多个 namespace,空值表示监听所有 namespace。它决定缓存观察的范围;Operator 身份能否读取和修改这些资源,还由 Kubernetes 的访问权限规则 RBAC 决定。

100 个 RayCluster 可以都在同一个 Kubernetes 集群中。这里的启动入口用一份 rest.Config 创建一个 Manager,这份客户端配置指定一套 Kubernetes API 地址和访问身份。增加 namespace 不会增加连接的 Kubernetes 集群;管理另一套 Kubernetes 集群,需要为它配置对应的 Operator 管理实例。

增加 Operator 副本还涉及选主:多个实例中由一个主实例负责协调,其余等待接管。KubeRay 默认开启 leader election,所以共享同一选举锁的副本,不能简单把协调并发数相加。若独立实例的观察范围重叠,各自队列的键锁也不能互相排斥;选举和接管在第五篇展开。

使用者同时修改对象时怎样处理

队列的串行规则作用于本 Controller 对同一个 RayCluster 的协调。与此同时,用户仍可以修改对象,其他 Controller 也可以写入相关资源。不同队列或不同对象的协调若涉及同一个下级资源,就需要由资源更新时的并发检查处理;进程内的共享内存则由相应的同步机制保护。

以一次状态更新为例:Controller 读到 A 的配置和状态,准备计算新的 status;在它提交之前,用户修改了 A 的 spec,API Server 已经保存了新版本。若直接用刚才读到的旧对象写回,就需要判断这个写入是否还基于有效版本。

Kubernetes 用 resourceVersion 标记对象版本。读取时客户端拿到这个标记,更新时原样带回,不自行推算或递增。普通带版本的 Update 若携带旧版本,API Server 会拒绝更新并返回 Conflict。这样,读取期间允许其他人修改,提交时再检查版本,这种方式称为乐观并发检查。

RayCluster 的状态写回函数 updateRayClusterStatus 使用 r.Status().Update(ctx, newInstance),把失败交回主流程。主流程在选择错误后执行下面的返回判断:

1
2
3
if err != nil || inconsistent {
return ctrl.Result{RequeueAfter: DefaultRequeueDuration}, err
}

这段代码可以同时返回错误和等待时间,实际重试规则由框架的 reconcileHandler 决定。普通错误非空时,通过 AddWithOpts 的 RateLimited 选项安排限速重试,并忽略 RequeueAfter;终止型错误有独立分支,不由本次错误自动重试。没有错误时,正的 RequeueAfter 才用于安排延时检查。

这次版本冲突后,本轮结束,后续协调重新读取、计算,再尝试写回。重新读取是其中必要的一步:同一份旧对象仍携带旧版本,原样重试还会冲突;下一轮能否写入,也取决于读到的版本是否已经更新。

版本检查以单次对象更新为单位。假设本轮先创建了一个 Pod,随后写 RayCluster status 时发生冲突,已创建的 Pod 会保留,下一轮结合配置、实际 Pod 和预期记录继续判断。这些跨对象操作分别生效,因此资源状态和 status 需要在后续协调中继续对齐。API 的 Patch 请求只描述部分修改,其冲突条件取决于具体写法,与这里的 Update 规则需要分别理解。

怎样判断是否需要提高并发

回到 100 个 RayCluster:全部稳定运行,与同时扩容、同时创建大量 Pod,给 Operator 带来的负载不同。稳定集群仍有周期性检查:RayCluster 主流程在收尾时安排下一轮,间隔可由环境变量配置;检查发现资源符合期望,就可以直接结束。

判断并发数是否合适,可以沿请求的处理过程看:等待处理的对象是否积压,协调协程是否长期占满,一轮处理是不是变慢,以及错误重试是否增加。controller-runtime 提供的指标分别对应这些位置:

要观察的情况 对应指标
已开始、尚未返回的协调是否长期占满名额 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 的集群数。

如果协调协程长期占满、许多不同对象持续排队,而单轮耗时和 API 请求表现稳定,提高并发值得验证。若积压主要来自一个正在处理的对象,仍受同键串行限制;若 Kubernetes API 已经变慢或限流等待增加,继续加协程还可能增加争用。

KubeRay 注册的 rest_client_request_duration_seconds 和 rest_client_rate_limiter_duration_seconds 分别帮助观察 API 请求耗时和客户端限流等待。客户端的 QPS、Burst 在 Manager 创建前配置,分别约束持续请求速率与突发额度,判断并发效果时也要结合这些限制。

小结:从请求安排到资源维护

协调协程主动取任务,队列后台分别整理请求、处理延时和分配任务;取得对象后,协调协程才执行 Reconcile。同一个集群的待处理请求可以合并,本轮执行期间也可以留下下一轮请求。待处理索引和执行锁定共同保证这些后续检查能够保留,同时避免同一对象的两轮协调重叠。

进入资源维护后,扩缩容预期衔接 API 写入与缓存观察之间的时间差,对象版本检查处理多个写入者的更新冲突。队列、预期记录和版本检查分别作用于请求安排、操作确认和资源写回,合在一起说明了多套集群的管理工作怎样推进。

评估并发配置时,需要结合每个 RayCluster 的 Pod 规模和变更负载,比较调整前后的排队耗时、错误速率及集群就绪耗时,才能判断这 100 个集群是否得到及时处理。下一篇把视角转向进程重启:资源已经创建,状态尚未写回时,Operator 怎样识别已有进展并继续协调。

参考资料


第 3 篇|KubeRay 并发:一个 Operator 如何管理 100 个 RayCluster
https://tanxinyu.work/kuberay-volcano-03-concurrency/
作者
谭新宇
发布于
2026年9月15日
许可协议