切换主题
调度与资源协调
问题:空闲节点为什么仍不能接任务
GPU 图像任务不能因为某台 CPU 节点负载低就派过去;版本不匹配的 Worker 也不能靠高权重成为候选。调度首先要满足资格,再在合格集合中比较偏好。
另一个容易混淆的问题是先处理哪个任务。队列优先级与老化决定任务顺序;节点打分决定某个任务去哪台机器。它们不是同一排序过程。
原理:四个阶段分别负责
- 队列选择:优先级、等待时长、公平性。
- 候选过滤:在线、状态、能力版本、硬标签、容量阈值。
- 偏好评分:标签、权重、已有负载。
- 原子预留:锁内复核资格与 slot,建立 Attempt 和命令。
纯函数适合做第二、三步,容易复现和测试;第四步必须与真实并发资源状态协调。两个调度器可能同时基于相同快照选同一 Worker。
最小例子:先过滤再排名
text
教学伪代码:
candidates = nodes.filter(online && versionMatches && memory >= required)
scores = candidates.map(preferredLabels + weight - runningLoad)
winner = sort(scores descending, stableID ascending).first()
在事务中复核 winner 的资格与剩余 slot
若被其他调度器占用,重新选择或保持排队稳定同分排序便于诊断,也可能让某个节点长期更受偏好;公平性是额外策略,不应被“确定性”自动替代。
交互 03
先准入,再比较偏好
| 候选 | 静态 CPU / 忙碌 Attempt | 资格与得分 |
|---|---|---|
| node-a | 4 核 / 2 | 1080000 |
| node-b | 8 核 / 0 | 150000 |
| node-c | 16 核 / 0 | 心跳过期 |
选择 node-a,得分 1080000。
简化至 CPU、单个优选标签和心跳,评分沿用当前案例:标签 × 1,000,000 + 权重 × 1,000 − 忙碌 Attempt × 10,000;同分按 Node ID。真实实现还包含其他资源、能力及事务内 slot 复核。
把 CPU 提升到 32 核会得到“无合格候选”。将心跳有效期改为 120 秒可让 node-c 入选;这是教学参数变化,不是推荐放宽生产离线判定。
方案比较:容量不等于利用率
静态容量表示机器能提供的资源,实时利用率表示此刻观测。利用率有延迟和噪声,仅靠“CPU 使用率低”不能保证同时启动多个重任务后不超载。
slot 限制并发个数,但一个轻任务与一个大内存任务可能同占一格。若需要按任务 CPU、内存和 GPU 隔离,要设计资源请求、原子预留、释放与实际运行限制。准入阈值不能假装是物理资源 reservation。
真实案例:确定性评分与事务复核
Select 过滤管理状态、Worker 运行和调度状态、心跳、精确能力版本、标签和静态资源;得分由优选标签、Node 权重和 RunningAttempts 构成。同分按 Node ID。
已核对的实现 · 本地代码快照
硬约束筛选与确定性评分 backend · c6dbe05c
server/internal/core/scheduling/domain/scheduler.go · 第 58–112 行
符号:Select · 核对日期 2026-10-02
来源与提交版本一致
// Select 先应用硬约束,再按得分降序、Node ID 升序确定唯一候选。
func Select(request Request, nodes []Node) (*Candidate, error) {
if request.Capability == "" || request.CapabilityVersion == "" {
return nil, errors.New("capability and version are required")
}
if request.Now.IsZero() || request.QueuedAt.IsZero() || request.HeartbeatTTL <= 0 {
return nil, errors.New("queue time, scheduling time, and heartbeat TTL are required")
}
candidates := make([]Candidate, 0, len(nodes))
for _, node := range nodes {
if request.NodeID != "" && node.ID != request.NodeID {
continue
}
if node.ManagementStatus != "ACTIVE" || node.Worker.RuntimeState != "READY" || node.Worker.SchedulingStatus != "ELIGIBLE" {
continue
}
if node.Worker.LastHeartbeat.Before(request.Now.Add(-request.HeartbeatTTL)) || !slices.Contains(node.Worker.CapabilityVersions, request.CapabilityVersion) {
continue
}
if len(QualificationReasons(node.TrustedLabels, node.Resources, request.RequiredLabels, request.MinimumResources)) > 0 {
continue
}
preferred := matchCount(node.TrustedLabels, request.PreferredLabels)
weight := node.Weight
if weight <= 0 {
weight = 100
}
score := int64(preferred*1_000_000 + weight*1_000 - node.RunningAttempts*10_000)
candidates = append(candidates, Candidate{
Node: node,
WorkerID: node.Worker.ID,
Score: score,
})
}
if len(candidates) == 0 {
return nil, errors.New("no eligible node")
}
slices.SortFunc(candidates, func(left, right Candidate) int {
if byScore := cmp.Compare(right.Score, left.Score); byScore != 0 {
return byScore
}
return cmp.Compare(left.Node.ID, right.Node.ID)
})
return &candidates[0], nil
}
// QualificationReasons 统一提供准入与只读诊断使用的资格过滤原因。
func QualificationReasons(labels map[string]string, capacity Resources, required map[string]string, minimum Resources) []string {
var reasons []string
if !matches(labels, required) {
reasons = append(reasons, "必选标签不匹配")
}
if capacity.CPUMillis < minimum.CPUMillis {
reasons = append(reasons, "CPU 静态规格不足")
}片段展示核对时的源码;完整文件指纹用于检测后续变化。这里的路径用于定位,不要求手机访问源码仓库。
调度事务与幂等 Attempt backend · c6dbe05c
server/internal/adapters/outbound/postgres/attempt_dispatch.go · 第 45–133 行
符号:dispatchAttempt · 核对日期 2026-10-02
来源与提交版本一致
func (dispatcher *AttemptDispatcher) dispatchAttempt(ctx context.Context, request executionports.CreateAttemptRequest, claim *pgqueries.SchedulingRequest) (*executionports.AttemptRef, error) {
stepID, err := uuid.Parse(request.StepID)
if err != nil {
return nil, fmt.Errorf("parse step ID: %w", err)
}
tx, err := dispatcher.database.BeginTx(ctx, pgx.TxOptions{
IsoLevel: pgx.ReadCommitted,
})
if err != nil {
return nil, err
}
defer func() { _ = tx.Rollback(ctx) }()
txQueries := dispatcher.queries.WithTx(tx)
lockKey := request.StepID + ":" + fmt.Sprint(request.AttemptNumber)
if err := txQueries.AcquireTransactionAdvisoryLock(ctx, lockKey); err != nil {
return nil, err
}
if _, err := txQueries.LockSchedulingRequest(ctx, pgqueries.LockSchedulingRequestParams{
ID: claim.ID,
ClaimID: claim.ClaimID,
}); err != nil {
return nil, err
}
executionIDValue, err := uuid.Parse(request.ExecutionID)
if err != nil {
return nil, fmt.Errorf("%w: execution ID", executiondomain.ErrDispatchInvalid)
}
execution, err := txQueries.LockExecutionByID(ctx, executionIDValue)
if err != nil {
return nil, err
}
if execution.Status != "CREATED" && execution.Status != "QUEUED" && execution.Status != "RUNNING" {
return nil, fmt.Errorf("%w: execution closed", executiondomain.ErrDispatchInvalid)
}
existing, err := txQueries.FindAttemptByStepAndNumber(ctx, pgqueries.FindAttemptByStepAndNumberParams{
StepID: stepID,
AttemptNumber: int32(request.AttemptNumber),
})
if err == nil {
if !existing.WorkerID.Valid {
return nil, errors.New("existing attempt has no worker")
}
if err := txQueries.CompleteSchedulingRequest(ctx, pgqueries.CompleteSchedulingRequestParams{
ID: claim.ID,
ClaimID: claim.ClaimID,
State: "DISPATCHED",
AvailableAt: timestamp(time.Now().UTC()),
}); err != nil {
return nil, err
}
if err := tx.Commit(ctx); err != nil {
return nil, err
}
existingWorker := uuid.UUID(existing.WorkerID.Bytes)
_ = dispatcher.waker.WakeWorker(existingWorker.String())
return &executionports.AttemptRef{
AttemptID: existing.ID.String(),
LeaseVersion: existing.LeaseVersion,
}, nil
}
if !errors.Is(err, pgx.ErrNoRows) {
return nil, err
}
now := time.Now().UTC()
workerID, generation, err := dispatcher.reserveWorker(ctx, txQueries, stepID, request, now)
if err != nil {
if errors.Is(err, errNoEligibleWorker) {
if commitErr := tx.Commit(ctx); commitErr != nil {
return nil, commitErr
}
}
return nil, err
}
contextRow, err := txQueries.FindStepExecutionAndWorkerNode(ctx, pgqueries.FindStepExecutionAndWorkerNodeParams{
StepID: stepID,
WorkerID: workerID,
})
if err != nil {
return nil, err
}
executionID, nodeID := contextRow.ExecutionID, contextRow.NodeID
if err := appendTimelineTx(ctx, txQueries, timelineRecord{
ExecutionID: executionID,
MilestoneKey: "step.worker.selected:" + stepID.String() + ":" + fmt.Sprint(request.AttemptNumber),
Stage: "SCHEDULING",
AttemptNumber: request.AttemptNumber,
State: "SUCCEEDED",
Component: "Scheduler",
Title: "已选定 Node 与 Capability 唯一 Worker",片段展示核对时的源码;完整文件指纹用于检测后续变化。这里的路径用于定位,不要求手机访问源码仓库。
当前纯算法中的动态利用率不参与决策,静态资源阈值是准入规则。具体派发适配器在事务里取得执行和调度锁,检查已有 Attempt,继续做 Worker 预留。不能据此声称每个任务拥有完整 CPU/GPU 硬隔离。
失败与边界:长期等待需要解释
无候选时记录原因比泛化成“系统错误”更有用。能力版本缺失、Worker 离线、资源不足是不同问题。节点清退或能力变化时,旧快照可能失效,派发需复核会话 generation 与清退状态。
资源释放要以 Attempt 事实及失效处理为依据,不能只凭 HTTP 断线立刻当作任务停止。等待时间与队列长度可用于发现饥饿和容量不足。
迁移练习与参考答案
练习:有两个 GPU 任务和一张 GPU,两个调度实例都看到卡空闲。如何避免双重分配?
参考答案:把 GPU 可分配单元建模并在同一受保护事务里预留,唯一约束或条件更新保证只有一个成功。纯打分在事务外运行后仍要复核;失败的请求重新排队。GPU 运行隔离与进程管理由执行环境落实,数据库标记不能独自提供物理隔离。