Skip to content

调度与资源协调 ​

问题:空闲节点为什么仍不能接任务 ​

GPU 图像任务不能因为某台 CPU 节点负载低就派过去;版本不匹配的 Worker 也不能靠高权重成为候选。调度首先要满足资格,再在合格集合中比较偏好。

另一个容易混淆的问题是先处理哪个任务。队列优先级与老化决定任务顺序;节点打分决定某个任务去哪台机器。它们不是同一排序过程。

原理:四个阶段分别负责 ​

  1. 队列选择:优先级、等待时长、公平性。
  2. 候选过滤:在线、状态、能力版本、硬标签、容量阈值。
  3. 偏好评分:标签、权重、已有负载。
  4. 原子预留:锁内复核资格与 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-a4 核 / 21080000
node-b8 核 / 0150000
node-c16 核 / 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 运行隔离与进程管理由执行环境落实,数据库标记不能独自提供物理隔离。

继续阅读:锁与租约、执行端恢复。

理解原理 · 分析取舍 · 用真实代码检验