XXL-Job、ElasticJob 等平台解决了"任务怎么跑、怎么分片"的问题(见 https://plumephp.com/distributed-job-scheduling/),但调度器内核的设计——任务模型、资源分配、优先级与抢占、延迟任务、故障恢复——才是决定大规模调度系统上限的关键。本文从零设计一个通用分布式调度器。
1. 调度器要解决什么问题
调度器(Scheduler)的职责是:在有限资源上,按约束与优先级,把待执行任务分配(schedule)到合适的执行节点,并保证最终全部执行完成。
┌─────────────┐ ┌─────────────┐
│ 任务提交 │ │ 资源节点 │
│ 优先级/依赖 │ │ CPU/内存/配额 │
└──────┬──────┘ └──────┬──────┘
▼ ▲
┌─────────────┐ │
│ 调度器核心 │───分配───►│
│ 队列/算法/抢占 │ │
└─────────────┘ │
│ │
▼ │
┌─────────────┐ │
│ 执行器/工作节点 │◄──执行────┘
│ 心跳/状态上报 │
└─────────────┘
1.1 调度器 vs 任务平台
| 维度 | 任务平台(XXL-Job 等) | 调度器内核(本文) |
|---|---|---|
| 关注点 | 任务管理、触发、分片、监控 | 资源分配、排队、抢占、恢复 |
| 输入 | 定时/手动任务 | 任意待调度任务(含资源描述) |
| 核心问题 | 什么时候执行 | 让谁在哪个节点用多少资源执行 |
| 代表 | XXL-Job / ElasticJob / PowerJob | Kubernetes 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/ 这类平台,也能支撑自研高吞吐调度系统的设计。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。