前几天我们内部一个叫“ax”的调度模块被新同事翻出来追问了好几次,起因是热词榜上突然挂了个“ax调度”,点进去发现大家说的其实是一类很朴素的问题:一堆异步任务挤在一起,到底怎么排、怎么跑、怎么在超时前收场。我仔细看了一下,发现这个“ax”就是我们两年前自研的异步执行调度核心,代码里还叫这个名。今天就把这套东西的来龙去脉、设计和踩坑记录完整地聊一遍,给正在折腾异步任务、批量请求、超时控制的同学一个能直接抄作业的参考。
所谓“ax调度”,本质上是把“任务提交”和“任务执行”彻底拆开,由调度器统一负责排队、优先级、超时、并发控制和重试。应用场景非常明确:上游接口不稳定、批量任务爆发、单次请求耗时不固定、系统不能被慢任务拖死。这套设计适合后端开发者、架构师,以及所有被“线程池被占满”“任务超时没人管”“突发流量打崩服务”折磨过的人。
1. 这个 ax 到底要解决什么问题:从业务痛点倒推调度设计
1.1 项目背景:乱序、超时与堆积的三重困境
最早我们要处理的业务是一批第三方数据源的批量采集任务,上游接口的响应时间从 50ms 到 30s 都有,偶尔还会直接挂掉。最初用线程池 + Future 硬扛,结果是线程池被慢请求占满,后续进来的任务全部排队等待,等了 2 秒、5 秒、10 秒,最后超时失败。而真正诡异的是,明明只调用了 200 个任务,线上却出现 2000 多个阻塞线程,连接池被打爆,CPU 飙到 90% 以上。
这个问题的本质是:线程是重量级资源,一个阻塞的线程不仅不干活,还占着内存和上下文切换的开销。而任务之间的依赖关系、优先级差异、超时后的补偿动作,在线程池模型里表达起来非常别扭。后来我们引入协程,但协程也不能无限开,否则内存照样爆,事件循环照样卡。于是就有了“ax调度”的需求:一个专门负责“决定谁先跑、跑多久、跑挂了怎么处理”的中间层。
1.2 三种主流方案的取舍比较
当时摆在面前的有三条路:第一种是继续用线程池,加大小并配置 CallerRuns 拒绝策略;第二种是引入消息队列,把所有任务丢到 MQ 里,消费者拉取执行;第三种就是自研一个像 ax 这样的协程调度器。逐个说结论:
- 线程池方案实现最简单,但超时控制仍然要依赖 Future.get(timeout),任务优先级不好做,而且线程阻塞是硬伤。
- MQ 方案解耦很漂亮,但延迟偏高,轻量级任务走 MQ 有点杀鸡用牛刀,而且还要额外维护一套 Broker。
- 协程调度器方案把执行单元切得足够小,由调度循环统一驱动,既能做优先级,又能精确控制超时和并发上限,还不需要额外组件。
最终我们选择了第三种,核心判断依据是:任务平均耗时短、数量大、突发性强、需要毫秒级响应,这种特征最适合在进程内做调度,而不是绕一圈走网络。
1.3 ax 调度的设计边界与量化目标
动手之前先定边界,否则很容易做成一个四不像的“万能调度框架”。我们给 ax 调度定了五个硬性目标:
- 并发上限可配置,默认不超过 500 个并发协程,防止突发流量打爆内存;
- 任务延迟可控,普通优先级任务从提交到开始执行不超过 100ms;
- 超时强制中断,单任务执行超过指定时间必须能释放资源,不依赖业务方自觉;
- 优先级可调整,线上紧急任务可以插队,但必须防止低优先级任务被饿死;
- 崩溃可观测,每个任务的生命周期都要有日志,便于事后复盘。
这五个目标直接决定了后面的数据结构选型和调度循环写法。
2. ax 调度的整体设计与分层思路
2.1 模块分层:提交、调度、执行三层隔离
ax 调度在架构上分了三个层次,每层只做一件事,职责边界非常清晰:
- 任务提交层:接收外部调用方提交的 Task,做参数校验、超时预算计算、生成全局唯一的 TaskID,并决定任务进入哪个队列。
- 调度核心层:维护多个优先级的待执行队列,由一个后台 goroutine 不断扫描,选出当前最该执行的任务,投递给执行器。顺便处理延迟任务的唤醒。
- 执行器层:真正干活的地方,从调度层拿到任务后启动子协程运行,同时挂上超时定时器,任务结束后回收资源并上报指标。
这个分层最大的好处是:外部逻辑永远不需要知道任务到底什么时候执行、由哪个协程执行,只要往提交层一丢,后续全由 ax 调度接管。我们后来接了十几个业务方,全部走同一套提交入口,没有任何一个业务方需要关心内部调度逻辑。
2.2 核心数据结构:优先级队列与时间轮的取舍
调度器的核心数据结构直接影响行为。先讲优先级:ax 用了三个优先级桶,分别是 high、normal、low,每个桶内部是一个 FIFO 队列。调度循环按“高优先先跑,但低优先级有最低配额”的策略取任务。这里没有用单一的大顶堆,因为我们不需要严格排序,只需要保证同优先级内先来先到,不同优先级间有配额,FIFO 队列 + 轮询配额的方式实现更简单、并发竞争更少。
再讲延迟队列,也就是需要“未来某个时间点再执行”的任务。我们一开始用的是 Go 的 time.Timer 逐个挂定时器,结果任务一多,定时器对象膨胀得厉害。后来换了时间轮:一个长度为 64 的环形数组,每个槽位存储该时刻需要唤醒的任务链表。调度循环每 tick 一次推进一格,把到期的任务重新塞回优先级队列。时间轮把定时器从 O(n) 的扫描变成 O(1) 的推进,实测在 5 万任务规模下 CPU 占用下降很明显。
2.3 参数设计与容量规划
参数不合理的调度器上线就是灾难。我们最初把队列容量设成无界,结果一次上游故障导致任务堆积了几百万,内存直接涨到 6GB,GC 频繁到服务不可用。后来所有队列都改成有界,并加了拒绝策略。
具体参数设计逻辑:
- 队列总容量 = 预估峰值 QPS × 单任务平均耗时 × 容忍堆积秒数。比如峰值 2000 QPS、平均耗时 500ms、容忍堆积 10 秒,那容量就是 2000 × 0.5 × 10 = 10000。
- 并发上限 = 目标吞吐量 × 单任务平均耗时,再加上 20% 的冗余。比如目标每秒完成 500 个任务,平均耗时 500ms,理论上需要 250 个并发,我们设 300。
- 超时预算 = 上游 P99 响应时间 × 1.5,不能拍脑袋设一个固定值,否则大量正常任务会被误杀。
3. 核心实现:手写一个可用的 ax 调度内核
3.1 任务模型与状态机
举个例子,ax 的任务结构长这样:
type Task struct { ID string Priority Priority Payload interface{} Timeout time.Duration MaxRetry int State TaskState }任务状态机很简单:Pending — Running — Succeeded / Failed / TimedOut / Retrying。状态转换的规则是:
- Pending 到 Running:调度循环将任务投递给执行器,同时启动超时定时器。
- Running 到 TimedOut:执行超时,调度器强制终止子协程,触发补偿回调。
- Running 到 Failed:业务执行返回 error,根据 MaxRetry 决定是 Retrying 还是 Failed。
- 任何状态到 Succeeded:任务结束,回收所有与任务相关的资源,包括上下文、定时器、traceID。
这里有个特别重要的细节:执行器必须为每个任务创建独立的子协程,而不是在调度协程里直接执行业务逻辑。否则一个任务阻塞就会把整个调度核心卡死。每个任务的超时定时器和取消函数都要在任务进入 Running 前注册,这样保证“任务开始就有兜底”。
3.2 调度循环的骨架代码
调度核心就是一个 for + select 循环,逻辑非常朴素:
func (s *Scheduler) loop() { ticker := time.NewTicker(time.Millisecond * 10) for { select { case <-s.stopCh: return case <-ticker.C: s.moveExpiredToReady() s.dispatch() case task := <-s.submitCh: s.enqueue(task) } } }dispatch 的核心是取任务和执行权的原子控制。伪代码如下:
func (s *Scheduler) dispatch() { if s.runningCount >= s.maxConcurrency { return } task := s.popTask() if task == nil { return } s.runningCount++ go func() { defer s.wg.Done() defer func() { s.runningCount-- }() ctx, cancel := context.WithTimeout(context.Background(), task.Timeout) defer cancel() done := make(chan error, 1) go func() { done <- task.Execute(ctx) }() select { case err := <-done: s.handleResult(task, err) case <-ctx.Done(): s.handleTimeout(task) } }() }这里要注意一个隐蔽的点:done 通道必须带缓冲。如果任务执行完恰好超时,主协程已经走到 ctx.Done() 分支,子协程写 done 就会阻塞,导致协程泄漏。带缓冲为 1 就能避免这个经典问题。
3.3 优先级与公平性的权衡
运行一段时间后发现一个现象:低优先级任务几乎永远得不到执行机会。高优先级任务在高峰期持续不断,low 队列里的任务被活活饿死。后来我们在 dispatch 里加了轮转配额:
- 每轮调度先取 high 队列,最多取 N 个,本轮 high 队列空了才轮到 normal。
- normal 队列每轮最多取 M 个,取完就强制转向 low。
- low 队列每轮至少取 1 个,哪怕 high 队列不为空,也要留出“最小执行窗口”。
这里 N 和 M 不能拍脑袋,我们的经验是 N 取并发上限的一半,M 取并发上限的三分之一。这样高优先级能保证快速响应,低优先级又不至于完全饿死。线上调整后,low 队列最长等待时间从原来的永远出不来降到 2 秒以内。
3.4 超时取消与优雅退出
超时取消不能真的把协程杀掉,而是通过 context 通知业务方自行让位。我们的 handleTimeout 逻辑是这样的:
func (s *Scheduler) handleTimeout(task *Task) { s.metrics.AddTimeout(task.ID) if task.OnTimeout != nil { task.OnTimeout(task) } if task.MaxRetry > 0 { task.MaxRetry-- task.State = TaskStateRetrying s.enqueueAfter(task, time.Second*time.Duration(s.retryBaseDelay)) } }业务方需要在自己的 Execute 函数里监听 ctx.Done(),比如:
select { case <-ctx.Done(): cleanup() return ctx.Err() case result := <-someResult: return result }最怕的是业务方完全不理会 context,超时后协程还在后台偷偷跑。我们的兜底方案是:执行器在超时后记录 goroutine 堆栈到日志,并且持续监控该任务的资源句柄,如果超过 3 秒还没释放,就上报告警。这个方案虽然不是物理杀协程,但至少会让问题暴露出来。
4. 实操记录:一次线上事故引发的参数调优
4.1 故障现场还原
上线 ax 调度一个月后,某次大促流量暴涨,监控系统报警:任务平均等待时间从 50ms 飙升到 12s,同时内存占用以每分钟 200MB 的速度上涨。查了 ax 调度的指标,发现 high 队列的长度从平时几百涨到几万,runningCount 一直顶在 500 没下来过,而单任务耗时从平均 300ms 涨到了 2s 以上。
原因链条很清楚:上游服务因流量过载变慢,单个任务耗时变长,在并发上限不变的情况下,单位时间能处理的任务数锐减,队列自然堆积。更麻烦的是,堆积的任务还在不断重试,重试又会把上游打得更慢,形成正反馈恶性循环。
4.2 排查过程和分析
首先确认是不是执行器代码有死锁,看了 pprof goroutine 堆栈,发现 90% 的 goroutine 都阻塞在上游 HTTP 调用上,说明不是死锁,纯粹是上游慢。然后看了重试计数,发现一个任务最高重试了 7 次,每次重试间隔只有 1 秒,相当于把上游当压力测试打。
我们当时的处理分三步:
- 立即动态调低并发上限到 200,同时调高任务超时时间,先止血;
- 关闭低优先级任务的重试,只保留 high 和 normal 的重试;
- 重试间隔从固定 1s 改为指数退避,第一次 1s、第二次 2s、第三次 4s,最大 30s。
4.3 参数计算的过程
大促期间预估任务峰值 QPS 是 3000,单任务 P99 耗时是 800ms,目标堆积时间不超过 5 秒。并发上限的计算:
3000 QPS × 0.8 秒 = 2400
加上 20% 冗余,理论上需要 2880 个并发协程。但这个数字太吓人了,我们评估了一下内存:每个任务携带的业务上下文约 2KB,每个 goroutine 初始栈 2KB,再加上调度器和队列开销,2880 并发非常接近单机的极限。最后做出的取舍是并发上限 600,靠队列兜住突发,同时把“超过峰值 1.5 倍的流量直接降级拒绝”,保证核心链路存活。
这里其实暴露了一个重要认知:并发上限不能只看单任务耗时,还要看任务的实际负载和资源占用。盲目把 maxConcurrency 调大,只是把问题从“等待时间长”变成“内存爆掉”。
4.4 调优后的效果与配置快照
调整后的核心配置供参考:
| 配置项 | 原值 | 新值 | 说明 |
|---|---|---|---|
| maxConcurrency | 500 | 600 | 提升吞吐但保留内存安全线 |
| high 队列容量 | 无界 | 10000 | 有界才可拒绝,避免 OOM |
| normal 队列容量 | 无界 | 30000 | 满足突发需求 |
| low 队列容量 | 无界 | 10000 | 低优先级任务允许丢弃 |
| 重试间隔 | 固定 1s | 指数退避 1s/2s/4s/8s/16s/30s | 避免重试风暴 |
| 超时时间 | 全局 2s | 按任务类别 0.5s/1s/3s | 不同业务差异化 |
上线后,任务平均等待时间降到 80ms,内存稳定在 1.2GB 左右,上游接口的可用率也恢复到了 99.9%。最明显的变化是,之前一直被慢任务拖死的连接池不再告警了。
5. 常见问题与避坑清单
5.1 任务饿死问题
现象:low 队列任务等待时间无限增长,业务方不断来投诉。 原因:调度循环永远优先取 high 队列,在持续高优先级压力下,低优先级没有机会执行。 解法:在 dispatch 里增加配额轮询,每轮强制给 low 队列一个最小执行窗口。配置项是 lowQueueGuaranteed,默认每轮至少 1 个,高峰时可以调大到 5。
5.2 超时回调泄漏问题
现象:任务已经超时,但 OnTimeout 回调里还在做一些耗时操作,导致超时流程越积越多。 原因:OnTimeout 是在调度核心协程里同步执行的,回调一慢,整个调度循环就卡住。 解法:OnTimeout 一律丢到独立协程池执行,不占用调度核心的时间片。更严格的做法是给 OnTimeout 本身也设置超时。这里要说一句:回调里的东西越少越好,最好只做标记和释放资源,别在里面写重逻辑。
5.3 队列堆积导致 OOM
现象:峰值流量下,队列容量设置得太大,任务对象占用内存过多,GC 频繁,最终 OOM。 原因:有界队列的容量设置没有和内存预算挂钩。 解法:在提交层做“当前队列长度 + 正在执行任务数”的总量控制,超过预算直接拒绝并返回错误,而不是让调用方无限排队。我们后来加了一个保护阈值:队列总量超过最大容量 80% 时,新任务直接走降级逻辑,不进入队列。
5.4 panic 导致调度器崩溃
现象:执行器里一个任务 panic,整个调度程序直接退出。 原因:执行器子协程没有捕获 panic,导致进程崩溃。 解法:每个任务执行体必须包裹 recover,并把 panic 信息记录到任务状态里。代码如下:
defer func() { if r := recover(); r != nil { fmt.Printf("task %s panic: %v\n", task.ID, r) s.handleResult(task, fmt.Errorf("panic: %v", r)) } }()资源回收也要在 defer 里做,不要用 defer 回收,却期望 panic 后还能继续走正常流程。
5.5 重试抖动问题
现象:任务失败后立即重试,仍然失败,然后再次重试,直接把下游打挂。 原因:重试策略没有考虑下游的恢复时间。 解法:重试间隔一定要带退避,而且要加随机抖动。计算公式是:
delay = min(maxRetryInterval, baseDelay × 2^retryCount) + random(0, jitter)
抖动的目的是防止多个任务同时重试,形成“同步雷击”效应。这个教训在大促事故里特别深刻,重试不加退避就是给上游送压力。
最后再分享一个实操心得
ax 调度这套东西整体实现下来,我最想强调的一点是:调度器只是骨架,真正决定系统稳不稳的是超时、重试和降级的配合。很多团队把调度器写成“一个能跑的 goroutine 池”就完事了,结果遇到故障时调度器反而成了帮凶:任务堆积、重试风暴、协程泄漏全来了。
如果你打算在自己的项目里实现类似 ax 的调度核心,我建议从最小的循环写起:先做一个固定并发上限的任务队列,跑通后再加优先级、加超时、加重试。每加一个功能都要配套对应的监控指标,至少要有队列长度、等待时间、超时次数、重试次数、执行耗时这五类。
另外补充一个小工具经验:验证调度器行为时,不要只靠单元测试,一定要写一个模拟慢任务的压测脚本。故意让一部分任务跑到超时边缘,看调度器的表现是否符合预期。很多 bug 都是在这种“半死不活”的任务压力下才能暴露出来。ax 调度后续还可以扩展批量聚合、分片分发、任务追踪,但核心的调度逻辑到现在我们都没动过,因为它足够简单,也足够稳。