Node.js异步任务调度器设计与实践:并发控制、优先级队列与重试策略
1. 项目概述与需求拆解1.1 “ax”到底是什么一个被团队内部叫了半年的小家伙先解释一下标题里这个“ax”。它是我和同事在去年一个后台项目里随手起的代号全称是async executor翻译过来就是“异步执行器”但大家习惯了直接叫它“ax调度”。简单说它是一套用 TypeScript 写的轻量级任务调度层跑在 Node.js 服务端负责管理所有出站 HTTP 请求和一批异步任务。它解决的核心问题只有三个并发请求把下游打崩、高优先级任务排队排到饿死、某个接口超时后其他请求跟着遭殃。当时为什么需要这种东西我们的服务要同时对接 20 多个第三方平台有的平台 QPS 限制很死比如每秒最多 5 个请求有的平台响应特别慢平均 3 秒起步。一开始大家各写各的A 同学用 Promise.all 批量发、B 同学在业务代码里自己数信号量、C 同学干脆所有请求都串行结果线上频繁出现两类事故一是瞬时并发冲上去后第三方封了我们的 IP二是某个慢接口把 Node.js 事件循环堵住整条请求链路全部超时。后来我们痛定思痛统一把请求收口到“ax”里由它统一调度情况立刻好转。1.2 明确要解决的三个痛点痛点清单是我们在项目初期一起梳理的建议你在动手写任何调度器之前也先把痛点列明白否则很容易把工具越做越复杂并发风暴某运营活动一上线后台瞬间进来几万条任务如果全部同时往下游发下游服务或者第三方接口直接被压垮。必须有一个窗口控制同时“在途”的请求数量。队列饿死常规队列都是先进先出但有些任务本身有时效性比如用户点了“刷新数据”之后 5 秒内需要返回结果如果前面排着几百个耗时的导出任务用户的请求会活活等死。需要有优先级机制。单点故障放大一个请求超时如果触发器设计得不好会连带取消一批无关任务甚至导致 worker 线程退出。调度器需要做隔离让单个失败不影响全局。1.3 哪些场景适合用它如果你现在的代码结构是“哪儿需要发请求就在哪儿 new 一个 Promise 直接发”那“ax”这类调度器非常适合你。尤其是下面几种情况对接多个外部 API且有独立配额限制。大量异步任务需要分批或并发执行执行前还要做参数校验、鉴权、限流。系统中有明确的优先级需求比如“人工触发任务”高于“定时批量任务”。想要把超时、重试、日志、监控全部统一收口而不是散落在各个调用方里。反过来也有不适用的情况如果你的请求总量本身很小一天就几百次那没必要上调度器直接 for 循环或者 Promise.all 就行。调度器是一层抽象抽象会带来学习成本和问题定位成本量小的时候不值得。2. 整体设计思路与核心原理2.1 为什么不能只用 Promise.all很多人的第一反应是控制并发用 Promise.all 不就行了吗Promise.all 只能处理“一次性批量发一批请求”它做不到动态调度。比如你要处理 1 万条消息一条一条地发Promise.all 会把 1 万个请求同时发出去下游瞬间崩溃。你当然可以自己写分批循环比如每 50 条一批用 Promise.all 跑完一批再跑下一批但这样做有不少问题批与批之间存在空窗期上一批全部完成后才开启下一批如果某一条特别慢整批都在等它吞吐量被最慢者拉低。没有优先级批与批之间是严格顺序的不是按任务紧急程度动态调整。没有失败重试和超时控制的统一入口你必须在业务代码里到处写 catch、写 AbortController。调度器的核心优势在于它维护一个“常驻执行池”只要有任务进来且当前在途任务数小于设定的并发上限就立刻启动。一个任务完成马上从队列里拉取下一个确保窗口始终是满的吞吐量接近理论最优。这个机制用一个简单的话说就是不要让机器闲着。2.2 调度器核心模型队列 限流器 工作池“ax”的设计参考了操作系统里的进程调度思路但在 JavaScript 里做了简化。核心模型由三部分组成任务队列Task Queue存储所有待执行的任务每个任务是一个对象包含执行函数、优先级、超时时间、重试次数、取消回调等元信息。信号量/限流器Semaphore / Limiter一个计数器记录当前在途任务数。每次任务开始时加 1任务结束时减 1。当计数达到上限时新任务只能进入等待队列。工作池Worker Pool可以理解成“固定数量的执行插槽”每个插槽在同一时刻只跑一个任务插槽空闲时自动从队列头部取任务。这个三件套并不是新鲜发明RabbitMQ、Kafka 的消费者组也有类似概念。但在前端/Node 生态里很多团队根本不用消息中间件只是内部 API 调用所以用这么轻的一层调度完全够用。条件允许的情况下你甚至可以把它封装成一个可注入的模块同时供 Web 服务端和定时任务脚本使用。2.3 优先级调度策略普通 FIFO 之外的取舍队列不能只做简单的先进先出FIFO因为真实业务里的任务天生就有优先级差异。我们在“ax”里实现的是加权优先级队列每个任务进入队列时带一个数字优先级数字越小优先级越高。每次取任务时从非空的最小优先级桶里取最早进入的任务。为防止低优先级任务被无限饿死设置一个“老化机制”任务在队列中每等待 10 秒优先级就临时上调一档最多上调到最高档。这个设计的判断依据是饿死一个后台任务意味着用户可能永远看不到一张报表饿死一个用户主动触发的刷新请求用户直接打客服电话。所以宁可偶尔“插队”也要保证所有任务最终都能被执行。实际测试中“老化机制”把低优先级任务的最大等待时间从原来的无限大降到了 45 秒以内。3. 核心细节解析与实操要点3.1 任务状态机设计别让任务处于“薛定谔状态”调度器里的任务不能用简单的“成功/失败”二进制状态来表示否则你很难排查问题。我们给任务定义了六个状态新建pending、排队中queued、执行中running、成功success、失败failed、已取消cancelled。这里有个经验一定要把“排队中”和“新建”分开。很多人的初始实现里没有“queued”状态任务一旦进来就直接标记为 running然后去抢并发插槽抢不到就放在数组里结果状态全乱套。实际上任务从进入队列那一刻起就应该进入一个明确的“等待调度”状态直到 worker 真正启动它才切换到 running。这样日志打出来你能一眼看出任务卡在了调度环节还是执行环节。任务对象上还需要挂一个attempts字段记录已重试次数。有一个常见的坑是把“重试”实现为一个新任务重新入队结果原任务还没结束新任务又进来了同一业务被重复执行。正确做法是在原任务对象上做重试更新 attempts 并重新入队而不是复制任务。3.2 并发上限、队列长度、超时时间怎么定这可能是最需要拍脑袋但也能部分计算的部分。我们最终把参数抽象成四个配置项参数我们的默认值调整思路concurrency建议 10~20根据下游接口能承受的 QPS 以及单请求平均耗时计算queueSize10000超过后新任务直接拒绝防止内存被打爆timeout15 秒略大于下游接口 P95 响应时间retries3 次重试太多会放大下游压力建议不超过 5先说concurrency怎么算。假设下游接口限制总 QPS 为 50单请求平均耗时 200ms那么单 worker 每秒可以完成 1/0.2 5 个请求。要达到 50 QPS理论上并发度 50 / 5 10。但这个公式里没有考虑响应时间波动所以我们一般再乘以 1.5 的安全系数最终设为 15。注意这不是让你一上来就调 15建议先设 1 或者 2用压测脚本慢慢往上加同时观察下游错误率。timeout的设计要小心不能拍脑袋设为“越小越好”。如果下游 P95 延迟是 8 秒你设 5 秒超时那么 5% 的正常请求会被误杀还会触发重试反而浪费更多并发窗口。我们的做法是看一整个星期内下游接口的延迟分布取 P95 值再加 2~3 秒余量。3.3 重试与退避策略别一失败就立刻再来一次重试不是“失败后重新执行”这么简单。如果你在第三方接口报 429 限流错误后立刻重试大概率还是 429而且会让对方更反感。我们的重试逻辑里引入了两种退避策略指数退避第 n 次重试前的等待时间 基础间隔 × 2^n同时加一个随机抖动0 ~ 50ms。错误类型感知对 HTTP 状态码做分类。429、5xx 可以重试 4xx 里的 400、401、403 不重试因为重试一百次也是同样的结果。这个分类至关重要。第一次实现时我们没区分错误类型所有失败都重试三次结果用户传了错误参数同一个错误请求被重复打到第三方三次第三方投诉我们刷接口。后来加入分类逻辑直接砍掉了大约 60% 的无意义重试。3.4 取消机制与请求中断用户取消后任务从队列彻底消失调度器必须支持取消单个任务和批量取消一组任务。这个在 Axios 场景中对应的是AbortController。我们在任务对象上存储一个abortSignal当收到取消指令时从等待队列中直接移除该任务如果还没开始执行。如果正在执行调用abortController.abort()触发底层请求中断。更新任务状态为 cancelled并尝试从“在途数量”里减 1。这里有个隐蔽的坑取消任务时如果底层请求已经发出并且收到了部分响应你不能直接假定资源已经释放。在 HTTP keep-alive 连接池里中断请求后连接可能被标记为坏连接如果不及时销毁连接池会被污染。我们的做法是在中断后强制关闭当前 socket 连接而不是等它自然回收。3.5 监控指标调度器不监控等于盲人骑瞎马调度器本身一定要暴露指标否则出了问题你根本定位不到是业务问题还是调度问题。我们为“ax”设计了一套最小指标集接入到了 Prometheusax_task_total累计接收任务数按优先级和类型打标签。ax_task_running当前在途任务数。ax_queue_depth当前队列深度。ax_task_wait_time任务从进队到开始执行的平均等待时间重点关注 P99。ax_task_failed_total失败任务数按错误类型分类。ax_worker_idle_ratioworker 空闲率如果长期为 0说明并发压满了需要扩 worker 或削峰。这套指标上线后帮我们抓到了一个大问题某业务方的请求一直优先级很高导致队列里 80% 都是低优先级任务它们平均等待时间达到了 7 分钟。如果没有监控这种“静默恶化”可能要出几次事故才能发现。4. 实操过程与核心实现4.1 最小可跑的调度器核心代码TypeScript下面给出一版简化但完整可用的调度器实现保留了核心机制并发控制、优先级、超时、重试和取消。完整代码在公司仓库里这里做了一些脱敏和剪裁但主干逻辑一致。type TaskPriority 1 | 2 | 3 | 4 | 5; interface TaskT unknown { id: string; priority: TaskPriority; execute: (signal: AbortSignal) PromiseT; timeoutMs: number; maxRetries: number; attempt: number; enqueueTime: number; cancellable: boolean; } interface SchedulerOptions { concurrency: number; defaultTimeoutMs?: number; defaultMaxRetries?: number; } class PriorityQueueT { private buckets: MapTaskPriority, TaskT[]; constructor() { this.buckets new Map(); } push(task: TaskT) { const list this.buckets.get(task.priority) || []; list.push(task); this.buckets.set(task.priority, list); } pop(): TaskT | undefined { for (let p 1; p 5; p) { const list this.buckets.get(p); if (list list.length 0) { return list.shift(); } } return undefined; } remove(taskId: string) { for (const list of this.buckets.values()) { const index list.findIndex((t) t.id taskId); if (index 0) { list.splice(index, 1); return true; } } return false; } size(): number { let total 0; for (const list of this.buckets.values()) total list.length; return total; } } export class AxScheduler { private queue: PriorityQueueunknown; private running: number; private workers: number; private options: RequiredPickerOptions; private abortControllers new Mapstring, AbortController(); constructor(options: SchedulerOptions) { this.queue new PriorityQueue(); this.running 0; this.options { defaultTimeoutMs: options.defaultTimeoutMs ?? 15_000, defaultMaxRetries: options.defaultMaxRetries ?? 3, concurrency: options.concurrency, }; this.workers options.concurrency; } submitT(task: { id: string; priority?: TaskPriority; execute: (signal: AbortSignal) PromiseT; timeoutMs?: number; maxRetries?: number; }): PromiseT { const wrapped: TaskT { id: task.id, priority: task.priority ?? 3, execute: task.execute, timeoutMs: task.timeoutMs ?? this.options.defaultTimeoutMs, maxRetries: task.maxRetries ?? this.options.defaultMaxRetries, attempt: 0, enqueueTime: Date.now(), cancellable: true, }; return new PromiseT((resolve, reject) { this.queue.push(wrapped); this.pump(); // 这里用队列任务完成回调的方式 resolve/reject简化版略 }); } private pump() { while (this.running this.workers) { const task this.queue.pop(); if (!task) return; this.runTask(task); } } private async runTaskT(task: TaskT) { this.running; const controller new AbortController(); this.abortControllers.set(task.id, controller); const timeoutTimer setTimeout(() { controller.abort(); }, task.timeoutMs); try { const result await task.execute(controller.signal); this.onTaskDone(task); return result; } catch (e) { task.attempt; if ( task.attempt task.maxRetries this.isRetryable(e) ) { const delay 100 * 2 ** (task.attempt - 1) Math.random() * 50; setTimeout(() { this.queue.push(task); this.pump(); }, delay); this.running--; return; } this.onTaskDone(task); throw e; } finally { clearTimeout(timeoutTimer); this.abortControllers.delete(task.id); this.running--; this.pump(); } } cancel(taskId: string) { if (this.queue.remove(taskId)) return true; const controller this.abortControllers.get(taskId); if (controller) { controller.abort(); return true; } return false; } }这段代码是可运行的但为了简洁省略了任务完成时的结果回传逻辑。实际使用时建议把submit内部封装成一个withExecutor辅助函数用onFulfilled、onRejected把Promise的 resolve/reject 和任务完成事件绑定起来避免任务在队列里重试时重复 resolve。4.2 和 Axios 结合中间件方式的接入技巧调度器本身不关心你执行什么它只调度“执行函数”。实际团队里我们主要用它来调度 Axios 请求。暴力做法是每个调用方都写scheduler.submit({ execute: () axios.get(...) })这样太累。我们做了两层封装第一层自己写一个scheduledAxios函数内部所有请求都走调度器import axios, { AxiosRequestConfig } from axios; let scheduler: AxScheduler; function initAxScheduler(options) { scheduler new AxScheduler(options); } function scheduledRequestT(config: AxiosRequestConfig): PromiseT { const taskId crypto.randomUUID(); return scheduler.submit({ id: taskId, priority: config.headers?.[x-priority] || 3, execute: (signal) axios.requestT({ ...config, signal }), }) as PromiseT; }第二层通过 Axios 自定义 adaptor 统一接管所有请求。这样业务代码几乎可以无感切换前提是你把项目中所有 axios 实例的默认adapter替换掉。我们当初是逐步灰度替换的先让一个内部服务试用再全量切换。axios.defaults.adapter function (config) { const taskId ${config.method}:${config.url}; return scheduler.submit({ id: taskId, priority: config.headers?.[x-priority] || 3, execute: (signal) { const normConfig { ...config, signal }; return axios.defaults.adapter?.(normConfig); }, }); };4.3 参数计算示例并发窗口从 3 到 30 的过程拿我们其中一个对接真实第三方支付接口的服务举例。该接口官方限制 QPS 为 30单请求平均耗时为 400msP95 为 900ms。我们按公式计算单 worker 每秒完成量 1 / 0.4 2.5 个请求。要达到 30 QPS理论并发度 30 / 2.5 12。考虑 P95 延迟偏大和网络抖动我们乘 1.5 安全系数得到 18取整后设置为 20。前三天我们设的是 15压测结果稳定后尝试调到 20。观察一周第三方接口没有出现 429我们的任务平均等待时间从 5 秒降到了 1.2 秒效果很明显。再往上调到 30 时出现了零星 429说明已经触碰到下游限流上限。最终我们保持 20这个例子说明并发度不是越大越好而是刚好卡在下游能承受的临界点之下。4.4 实测数据队列积压从分钟级降到秒级改造前后的数据对比很直观。改造前我们用最原始的“每 100 个请求一批”循环方式积压 5000 个任务时最后一批任务的平均等待时间接近 4 分钟。改造后同样 5000 个任务、并发度 20任务平均等待时间 23 秒P95 等待时间 41 秒。因为我们把那些“用户主动请求”设为优先级 1 之后它们的平均等待时间不到 3 秒同时低优先级任务的等待时间也没有超过 60 秒老化机制起了作用。5. 常见问题与排查技巧实录5.1 死锁和饥饿优先级队列的隐藏陷阱有一次我们上线后接到反馈某个接口的请求全部超时。查日志发现ax_queue_depth从几百涨到了几千但ax_task_running一直等于 0。这说明根本没有任务在执行调度器卡死了。最后定位到原因某个高优先级任务在execute里抛出的错误没有被 catch导致 worker 的running计数在异常传播过程中没有正确减掉从而整个 worker 永久“占着茅坑不拉屎”。排查技巧是写调度器时pump()方法必须在 finally 块里反复调用并且running计数不能放在 try 块内自增。最简单的防御性写法是 RanTask 入口加一个独立变量的状态校验每次退出时都确认 running 数量与真实活动任务相等。饥饿问题则不太一样它是低优先级任务长时间不被执行。我们用老化机制解决了但老化参数不能太激进否则高低优先级会混成一团。建议老化间隔设为平均响应时间的 5~10 倍我们设的 10 秒。5.2 请求没有真正取消连接池污染上面提到的连接池污染我再展开说一下。刚开始实现取消功能时我们只调了controller.abort()结果发现大量请求还是会在第三方平台产生日志而且过了一段时间后本地 socket 文件描述符数量飙升。原因就是连接池中那些被中断的连接没有被正确关闭Node.js 认为连接还活着继续复用但实际上远端已经不可用。解决办法是给 Axios 实例设置httpAgent: false或者使用专门的 agent 配置在取消时销毁该连接。如果是自己管理 socket一定要在 finally 里判断当前连接的状态如果closed为 false 就主动销毁。5.3 并发数调大反而更慢别忽视下游连接数和线程池有一个反直觉的现象并发数从 10 调到 30 后吞吐量没有线性上升反而下降了。原因很简单下游服务是一个 Java 应用它的 Tomcat 线程池最大线程数是 20。我们 30 个并发请求同时打过去20 个线程在干活剩下 10 个请求在 Tomcat 队列里排队而 Tomcat 的默认队列策略会对这些请求产生额外延迟。这个问题的判断标准是观察调度器的平均任务耗时如果并发调大后平均耗时明显上升说明下游已经到瓶颈了。此时应该把并发数调回临界值而不是继续堆并发。从这个例子你应该能体会调度器的并发上限本质上是一个“双方约定”的值必须结合下游能力来定。5.4 队列积压后的快速清空策略突发事件导致队列积压几千个任务时如果按正常速度处理可能要好几分钟用户等不了。我们的快速清空策略分三步先把新任务拒之门外将 queueSize 临时设为 0 或直接返回 503。将低优先级任务的调度权重临时降低确保高优先级任务以最快速度处理。动态提升并发数但绝对不能一次性加满而是每次加 5 个观察下游错误率如果 429 变多就回退。这套策略在两次大促活动中成功保住了用户体验我们的监控里能看到瞬时并发拉高后 1 分钟内队列恢复到正常水位。5.5 不要为了清空队列而重试失败任务另一个反面教训当时我们为了快速清空积压遇到失败就立刻重试结果某个下游已经完全不可用重试只是在白白消耗资源最终队列没清空反而又堆了一堆重试任务。后来我们加了“熔断”逻辑如果同一个下游连续失败超过 10 次直接将该下游的所有请求标记为快速失败间隔 30 秒后再开始试探。正常业务下低于 10 次的连续失败几乎是不会出现的所以这个阈值不会误伤。6. 写在最后的一点经验“ax”这个调度器我们迭代了快半年最大的收获不是那几千行代码而是搞清楚了一件事大部分后端系统的瓶颈不在代码性能而在资源协作。谁先执行、谁后执行、谁可以放弃、谁必须重试这些策略比用多牛的语言和框架重要得多。如果你也要在团队里做类似的调度层我的建议是不要一上来就引入 Redis 队列、消息中间件这些重武器先用一个内存里的优先级队列跑通流程把指标埋清楚等真的出现跨进程调度需求时再演进。过程中要记得把每个参数的调整原因写成文档不然三个月后没人记得 concurrency 为什么是 15。最后分享一个小技巧调度器上线头一个月每天拉一次 P99 等待时间和队列深度做日报它会帮你提前发现很多奇怪但真实的流量特征。