分布式调度器设计

分布式调度器设计:任务模型、资源管理、优先级与抢占、延迟任务与重试、调度器高可用

XXL-Job、ElasticJob 等平台解决了"任务怎么跑、怎么分片"的问题(见 https://plumephp.com/distributed-job-scheduling/),但调度器内核的设计——任务模型、资源分配、优先级与抢占、延迟任务、故障恢复——才是决定大规模调度系统上限的关键。本文从零设计一个通用分布式调度器。

1. 调度器要解决什么问题

调度器(Scheduler)的职责是:在有限资源上,按约束与优先级,把待执行任务分配(schedule)到合适的执行节点,并保证最终全部执行完成。

        ┌─────────────┐       ┌─────────────┐
        │  任务提交     │       │  资源节点     │
        │  优先级/依赖   │       │  CPU/内存/配额 │
        └──────┬──────┘       └──────┬──────┘
               ▼                     ▲
        ┌─────────────┐              │
        │  调度器核心    │───分配───►│
        │ 队列/算法/抢占 │              │
        └─────────────┘              │
               │                     │
               ▼                     │
        ┌─────────────┐              │
        │ 执行器/工作节点 │◄──执行────┘
        │ 心跳/状态上报  │
        └─────────────┘

1.1 调度器 vs 任务平台

维度任务平台(XXL-Job 等)调度器内核(本文)
关注点任务管理、触发、分片、监控资源分配、排队、抢占、恢复
输入定时/手动任务任意待调度任务(含资源描述)
核心问题什么时候执行让谁在哪个节点用多少资源执行
代表XXL-Job / ElasticJob / PowerJobKubernetes Scheduler / Mesos / 自研

1.2 调度的本质权衡

  • 公平性:多个任务源/租户之间资源不能失衡
  • 效率:资源利用率高,任务等待短
  • 优先级:高优任务优先获得资源
  • 确定性:相同输入得到可预期调度结果

调度器设计就是在这四者之间找平衡。

2. 任务模型

2.1 任务描述

一个调度任务应包含资源需求、优先级、依赖与约束:

type Task struct {
    ID            string            `json:"id"`
    Queue         string            `json:"queue"`          // 所属队列
    Priority      int               `json:"priority"`       // 数值越大优先级越高
    CPUReq        int64             `json:"cpuReq"`         // 所需 CPU 毫核
    MemReq        int64             `json:"memReq"`         // 所需内存 MB
    Dependencies  []string          `json:"dependencies"`   // 依赖任务 ID
    Deadline      int64             `json:"deadline"`       // 截止时间戳
    Affinity      map[string]string `json:"affinity"`       // 节点亲和(标签)
    MaxRetries    int               `json:"maxRetries"`
    Timeout       int64             `json:"timeout"`        // 单次执行超时
    Payload       string            `json:"payload"`
    State         TaskState         `json:"state"`
}

2.2 任务状态机

PENDING ──► QUEUED ──► RUNNING ──► SUCCEEDED
    │          │          │
    │          │          ├──► FAILED ──► 重试 ──► QUEUED(若未超次数)
    │          │          │
    │          │          └──► TIMEOUT ──► 重试 / 判定失败
    │          ▼
    └────► CANCELED / REJECTED(资源不足且不可等)

2.3 DAG 依赖

有依赖的任务用 DAG 表达,调度器先调度入度为零的任务,完成后推进下游:

func (d *DAG) ReadyTasks() []string {
    ready := []string{}
    for id, node := range d.Nodes {
        if node.InDegree == 0 && !node.Done && node.Running == 0 {
            ready = append(ready, id)
        }
    }
    return ready
}

3. 资源管理

3.1 资源抽象

  • 可量化的资源:CPU(毫核)、内存(MB)、磁盘、GPU
  • 不可量化但需约束的:连接数、带宽、文件句柄
  • 节点容量:每个工作节点上报可用资源,调度器维护集群资源视图
集群资源视图:
  node-a: { cpu: 8000m, mem: 32Gi, used_cpu: 2000m, used_mem: 8Gi, free_cpu: 6000m, free_mem: 24Gi }
  node-b: { cpu: 4000m, mem: 16Gi, used_cpu: 4000m, used_mem: 16Gi, free_cpu: 0m,   free_mem: 0Gi  }
  node-c: { cpu: 8000m, mem: 32Gi, used_cpu: 1000m, used_mem: 4Gi, free_cpu: 7000m, free_mem: 28Gi }

3.2 调度算法

算法原理特点适用
FIFO先来先服务简单、无抢占低负载、简单批处理
加权公平队列按队列权重分配资源公平、防饿死多租户、多队列
最少负载挑剩余资源最多的节点均衡、实现简单通用
Bin Packing尽量填满节点高利用率、节省成本资源密集场景
优先级抢占高优任务抢占低优任务资源高优优先、复杂在线/离线混合

3.3 加权公平队列实现

public class WeightedFairQueue {
    // 每个队列维护一个虚拟运行时间,调度时选虚拟时间最小的队列出队
    private final Map<String, QueueState> queues = new ConcurrentHashMap<>();

    public synchronized Task scheduleNext() {
        String winner = null;
        long minVirtual = Long.MAX_VALUE;
        for (var e : queues.entrySet()) {
            if (!e.getValue().isEmpty()) {
                long v = e.getValue().virtualTime / e.getValue().weight;
                if (v < minVirtual) {
                    minVirtual = v;
                    winner = e.getKey();
                }
            }
        }
        if (winner == null) return null;
        Task t = queues.get(winner).dequeue();
        queues.get(winner).virtualTime += 100; // 每出队一个任务累加
        return t;
    }
}

4. 优先级与抢占

4.1 优先级队列

多个优先级队列 + 严格/权重混合策略:

优先级 0(紧急):立即调度
优先级 1(重要):等待时间可短
优先级 2(普通):按 FIFO
优先级 3(后台):低峰执行 / 可被抢占

调度顺序:先看高优先级队列是否为空,再逐级向下;同一优先级内按 FIFO。

4.2 抢占(Preemption)

当高优任务到达但资源不足时,可以抢占低优任务的资源。抢占设计要点:

设计点方案
抢占对象抢占同队列中最低优先级的 RUNNING 任务
优雅退出给被抢占任务宽限期(grace period),允许其保存状态
抢占成本考虑任务已运行时长(快完成的任务不抢)
防震荡被抢占任务不立即重排队,退避后再试
func (s *Scheduler) preempt(high *Task) bool {
    victims := s.findVictims(high)      // 找低优 RUNNING 任务
    sort.Slice(victims, func(i, j int) bool {
        if victims[i].Priority != victims[j].Priority {
            return victims[i].Priority < victims[j].Priority
        }
        // 快完成的任务优先保留:剩余运行时长排序
        return victims[i].RemainingEst > victims[j].RemainingEst
    })
    freed := int64(0)
    for _, v := range victims {
        s.kill(v, "preempted by "+high.ID)   // 发起优雅退出
        freed += v.CPUReq + v.MemReq
        if freed >= high.CPUReq+high.MemReq {
            return true
        }
    }
    return false
}

4.3 优先级反转与继承

高优任务可能被低优任务持有的锁阻塞(优先级反转)。简单做法:持有锁的任务临时继承等待者的优先级(优先级继承),减少高优任务等待。

5. 延迟任务与重试

5.1 延迟队列实现

延迟任务(延时执行)用时间轮(Timing Wheel)或堆实现:

// 最小堆实现延迟队列
public class DelayQueue {
    private final PriorityQueue<DelayedTask> heap = new PriorityQueue<>(
        Comparator.comparingLong(DelayedTask::getDeadline));

    public void add(DelayedTask task) {
        heap.offer(task);
    }

    public List<DelayedTask> pollExpired(long now) {
        List<DelayedTask> ready = new ArrayList<>();
        while (!heap.isEmpty() && heap.peek().getDeadline() <= now) {
            ready.add(heap.poll());
        }
        return ready;
    }
}

时间轮把到期任务按"槽位"组织,复杂度 O(1),适合海量定时/延迟任务。到期任务从时间轮移动到就绪队列交给调度器。

5.2 重试与退避

任务失败后重试,必须控制重试风暴:

退避策略行为适用
固定间隔每次等待相同时间简单场景
指数退避2^n 倍递增网络类故障
指数退避 + 抖动退避基础上加随机抖动,防止 thundering herd分布式系统首选
封顶超过最大间隔后固定防止无界等待
func nextRetryDelay(attempt int, maxBackoff time.Duration) time.Duration {
    exp := time.Duration(1<<uint(min(attempt, 10))) * time.Second
    if exp > maxBackoff {
        exp = maxBackoff
    }
    // 加抖动,避免重试风暴
    jitter := time.Duration(rand.Int63n(int64(exp / 4)))
    return exp/2 + jitter
}

5.3 重试与幂等

重试必须配合任务的幂等执行,否则重试会放大副作用。任务执行器应按 task_id 去重(参考 https://plumephp.com/distributed-idempotency-reliability/ 的幂等设计)。

6. 调度器高可用

6.1 主备选举

调度器自身必须是高可用的,否则单点故障会停摆整个调度。经典方案:

调度器实例 A(Leader):负责调度决策
调度器实例 B(Follower):待命
调度器实例 C(Follower):待命

Leader 通过分布式锁/租约选主(etcd/Redis/ZooKeeper)
Leader 宕机 → 租约过期 → B 竞选为新 Leader

选主机制可复用 https://plumephp.com/distributed-locking/ 与 https://plumephp.com/zookeeper-coordination/ 中的实现。

6.2 状态持久化与恢复

调度器的队列、分配关系、任务状态必须持久化,Leader 挂掉后新 Leader 能从持久化状态恢复:

持久化存储(etcd / MySQL):
  pending_queue:待调度任务
  running_tasks:正在运行的任务 → 节点、开始时间
  node_registry:节点资源视图
  task_state:任务状态机进度

恢复流程:
  1. 新 Leader 选主成功
  2. 加载持久化状态,重建队列与资源视图
  3. 对所有 RUNNING 任务做"重新认领"(leader 变更导致执行状态不确定)
  4. 认领超时的任务重新入队

6.3 节点故障处理

工作节点宕机时,调度器需重新调度其上的任务:

func (s *Scheduler) OnNodeDown(nodeID string) {
    for _, t := range s.runningOnNode(nodeID) {
        // 依据任务类型决定:可重跑 → 重新入队;不可重跑 → 标记 FAILED 告警
        if t.Restartable {
            s.enqueue(t, WithRetryReason("node down: "+nodeID))
        } else {
            s.markFailed(t, "node down, non-restartable")
        }
    }
    // 更新资源视图,把节点标记不可用
    s.nodes.MarkDown(nodeID)
}

7. 调度器实现示例

7.1 核心调度循环

type Scheduler struct {
    queue   *WeightedFairQueue
    nodes   *ResourceManager
    store   *TaskStore
    delay   *TimeWheel
}

func (s *Scheduler) Run(ctx context.Context) {
    for {
        select {
        case <-ctx.Done():
            return
        case <-time.After(50 * time.Millisecond): // 调度节拍
            s.tick()
        }
    }
}

func (s *Scheduler) tick() {
    now := time.Now()
    // 1) 时间轮到期任务进入就绪队列
    for _, t := range s.delay.PollExpired(now.UnixMilli()) {
        s.queue.Enqueue(t)
    }
    // 2) 尝试调度就绪任务
    for {
        task := s.queue.Next()
        if task == nil {
            break
        }
        node := s.nodes.BestFit(task)
        if node == nil {
            // 资源不足:尝试抢占高优前的低优任务
            if task.Priority >= PreemptThreshold && s.preempt(task) {
                continue
            }
            s.queue.RequeueBack(task) // 放回队首等待
            break
        }
        s.dispatch(task, node)
    }
}

7.2 分布式一致性要点

  • 选主:用 etcd 租约 + 序号保证只有一个 Leader
  • 双活避免:调度决策写入持久化存储时用事务/条件更新,防止双 Leader 重复调度
  • 分配幂等:每个任务有唯一调度记录 (task_id, node_id, attempt),重复分配直接忽略

总结

模块关键设计要点
任务模型资源需求 + 优先级 + DAG 依赖 + 状态机状态可恢复、依赖可推进
资源管理节点资源视图 + 调度算法公平/效率/优先级权衡
优先级与抢占多级队列 + 优雅抢占 + 优先级继承防震荡、防反转
延迟与重试时间轮 + 指数退避抖动延迟 O(1),重试防风暴
高可用选主 + 状态持久化 + 节点故障重调度可恢复、不双活

分布式调度器是任务平台与编排系统(如 Kubernetes Scheduler)的底层内核。理解它的任务模型、资源分配、抢占与恢复机制,既有助于用好 https://plumephp.com/distributed-job-scheduling/ 这类平台,也能支撑自研高吞吐调度系统的设计。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「distributed-systems」更多文章

  1. 分布式数据库前沿深度解析:TiDB、Spanner 与 CockroachDB 的共识与事务实现
  2. 异地多活与容灾架构深度解析:同城双活、两地三中心与多活设计
  3. 幂等设计与消息可靠性:不丢不重、防止重复消费的分布式基石