
前几天我们内部一个叫“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 参数设计与容量规划参数不合理的调度器上线就是灾难。我们最初把队列容量设成无界结果一次上游故障导致任务堆积了几百万内存直接涨到 6GBGC 频繁到服务不可用。后来所有队列都改成有界并加了拒绝策略。具体参数设计逻辑队列总容量 预估峰值 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 调优后的效果与配置快照调整后的核心配置供参考配置项原值新值说明maxConcurrency500600提升吞吐但保留内存安全线high 队列容量无界10000有界才可拒绝避免 OOMnormal 队列容量无界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 调度后续还可以扩展批量聚合、分片分发、任务追踪但核心的调度逻辑到现在我们都没动过因为它足够简单也足够稳。