数据截至 (上游 commit 9194847b2cf8)
调度器:请求怎么找到一台机器
这一章讲 Beam 控制面的心脏:一个
ContainerRequest(「请给我跑一个这样的容器」)是怎么被匹配到一台具体 worker 的,匹配不上又怎么办。代码主要在pkg/scheduler/。
1. 它要解决的小问题
抽象层会源源不断地说「我要再开 3 个容器」。但容器不能凭空出现——得有一台有足够 CPU/内存/对的 GPU 型号、且没被占满的机器来跑它。调度器就是那个「把请求和机器撮合」的中间人;撮合不上时,它还得去开新机器。
2. 思路 / 直觉:一条 Redis backlog + 一个批处理循环
Beam 没有用「请求直达 worker」的同步模型,而是经 Redis 解耦:
- 抽象层调
Scheduler.Run(request),请求被推进 Redis 的 backlog(一个队列)。 - 调度器另有一个独立循环
StartProcessingRequests,不停从 backlog 批量弹出请求来处理。
这样做的好处:生产(要容器)和消费(撮合)速度解耦,且调度状态天然可以多副本共享 backlog。
Scheduler.Run(req) StartProcessingRequests()(独立 goroutine)
│ │ 每 50ms
│ 设置并发配额 ▼
│ 检查重复容器 PopN(512) 一批请求
▼ │
addRequestToBacklog ──RPush──► Redis backlog ──LPop──► 逐个 processRequest
关键常量(scheduler.go:25-31):批大小 requestProcessingBatchSize=512、轮询间隔 requestProcessingInterval=50ms。
3. 入口:Scheduler.Run 做了什么
Run 不直接调度,它只做「入队前的准备」(scheduler.go:212 func (s *Scheduler) Run):
- 查重:若该
ContainerId已是 pending/running,返回ContainerAlreadyScheduledError,避免重复起。 - 附检查点:若开了自动检查点且有可用检查点,把
request.Checkpoint填上(为后面 CRIU 恢复铺路)。 - 并发配额:
getConcurrencyLimit取 workspace 的并发上限并SetContainerStateWithConcurrencyLimit落库。 - 入 backlog:
addRequestToBacklog。
注意一个细节:addRequestToBacklog (scheduler.go:990) 里第一次(RetryCount==0)是立即入队、不延迟;之后每次重试都走指数退避 calculateBackoffDelay(1s 起、5s 封顶,scheduler.go:1100),并在超过 maxScheduleRetryCount=120 次或 maxScheduleRetryDuration=10min 后放弃。
4. 核心:三级调度顺序
请求被弹出后交给 schedulingAttempt.run()(reserve.go:29)。这里把调度顺序写得非常明确:
run():
① scheduleOnAvailableWorker() 就绪容量优先:有空闲 worker 直接派
└ 成功 → 返回
② reservePendingWorkerCapacity() 其次:预留一台「正在启动中」的 worker 的容量
└ 成功 → 重新入队等它起来
③ tryPrivatePoolFallback() 私有池兜底
④ provisionWorker() 最后:调云厂商 API 开一台新机器
源码注释原话(reserve.go:35-36):「use ready capacity first, wait on pending capacity second, and only provision when neither path can fit」。这条「能用现成的就别开新机」的优先级,是省钱省时间的核心。
4.1 第一级:选一个空闲 worker(过滤 + 打分)
selectWorkerFromWorkersByStatus(scheduler.go:869)是撮合的核心,分两步。
第一步:一串过滤器,逐层筛掉不合格的 worker。 这是并列的几道关卡,用表更清楚:
| 过滤器 | 函数 | 筛掉谁 |
|---|---|---|
| 池选择器 | filterWorkersByPoolSelector | 请求指定了私有池 → 只留该池的 worker;没指定 → 排除「要求选择器」的 worker |
| 私有 agent 存活 | filterLivePrivateAgentWorkers | 私有池里排除掉「机器已离线」的 worker |
| 资源 | filterWorkersByResources | CPU/内存不够、GPU 型号不匹配、被 cordon 的 |
| 标志 | filterWorkersByFlags | 非抢占请求排除掉可抢占 worker |
| 状态 | filterWorkersByStatus | 只留指定生命周期状态(如 Available) |
资源过滤里有个巧妙的内存超额预留:capacityMemoryForScheduling(scheduler.go:809)把请求内存乘 1.25 (memory*125+99)/100——给容 器留 25% 余量,避免刚好卡死触发 OOM。
第二步:给活下来的 worker 打分,取最高分。
# 示意,非源码:评分逻辑对应 scoreWorkerForRequest
def score(worker, request):
s = worker.priority # 池优先级是基础分
if worker.status == "available":
s += 10 # 现成可用的 +10(scoreAvailableWorker)
if request.requires_gpu():
s -= gpu_priority_modifier(...) # GPU 型号在请求偏好列表里越靠后,扣得越多
return s
# 重点看:分相同时,再比「空闲容量」——GPU 数权重极高(×1_000_000)
真实实现 scoreWorkerForRequest(scheduler.go:949)+ 平局打破 workerFreeCapacityScore(scheduler.go:975)。后者把空闲 GPU 数乘以 100 万,意思是优先把活塞到 GPU 最空的机器,让 GPU 利用率更均衡。
4.2 第二级:预留「待启动」worker 的容量
如果没有现成空闲 worker,但已经有一台机器正在启动中(status=Pending),调度器不会傻等也不会立刻又开一台,而是在那台未来的机器上先预扣容量——provisioningTracker.reserveCapacity(reserve.go:350)。
这块用一个 provisioningTracker 维护两张映射:reservations(worker → 预留)和 requestReservations(请求 → 预留)。预留成功后请求重新入队,等那台机器起来就能直接命中第一级。预留有 TTL 兜底(pendingWorkerReservationTTL=30s)防止泄漏。