引言
一个 HTTP 请求能承受的工作量是有上限的:发送邮件、生成报表、转码视频、调用第三方 API,这些操作耗时从数百毫秒到数分钟不等。把它们塞进请求-响应周期里,用户要等,超时还会导致重复执行。后台任务队列就是把这些工作从请求路径上剥离出去的标准手段。
在 Node/TypeScript 生态里,BullMQ 是使用最广的队列实现之一:它基于 Redis,提供任务重试、延时、优先级、限流与并发控制,并有完善的类型定义。但它也带来一批新的工程问题——重试语义、幂等性、去重、死信、优雅关闭,这些才是线上真正会出故障的地方。
本文聚焦 BullMQ 在 TypeScript 中的落地:从核心模型讲起,覆盖 payload 泛型、重试退避、幂等去重、定时延时、死信告警、并发限流,最后给出可观测与优雅关闭实践。
目录
- 1. 为什么需要任务队列
- 2. BullMQ 核心模型
- 3. 任务类型与 payload 泛型
- 4. 重试与退避策略
- 5. 幂等与任务去重
- 6. 定时与延时任务
- 7. 失败告警与死信队列
- 8. Worker 并发与限流
- 9. 可观测与任务追踪
- 10. 优雅关闭与生产实践
1. 为什么需要任务队列
1.1 同步请求的边界
请求-响应周期适合「快、可预期、可重试成本低」的操作。一旦操作耗时超过几百毫秒、依赖外部服务、或可能失败需要重试,就该移出请求路径。否则用户等待变长、网关超时、失败无法补偿。
1.2 队列带来的能力
| 能力 | 含义 | 典型场景 |
|---|---|---|
| 异步化 | 立即返回,后台执行 | 发邮件、生成报表 |
| 削峰 | 缓冲突发流量 | 秒杀后置处理 |
| 重试 | 失败自动重放 | 调用第三方 API |
| 延时 | 到点才执行 | 订单超时取消 |
| 限流 | 控制执行速率 | 受限配额接口 |
1.3 方案选型
BullMQ 基于 Redis,功能完整、生态成熟、类型好;Bee-Queue 更轻量但功能少;pg-boss 基于 PostgreSQL,适合不想引入 Redis 的团队;云上则可用 SQS 等托管服务。若已有 Redis,BullMQ 通常是默认选择。
1.4 引入队列的代价
队列不是免费的:多一个需要运维的中间件、任务状态需要额外存储、调试从「看堆栈」变成「看队列面板」、失败语义与幂等要自己设计。这些代价在引入前就该被评估。
一句话总结:把耗时、易失败、可延后的工作移出请求路径——队列带来异步、削峰、重试与限流能力,但也引入了新的运维与幂等设计成本。
2. BullMQ 核心模型
2.1 三个核心概念
Queue 是任务的生产与查询入口;Worker 从队列取任务并执行;Job 是单条任务,携带名字、数据、状态与重试信息。QueueEvents 则订阅任务生命周期的全局事件,用于告警与追踪。
2.2 最小示例
import { Queue, Worker } from "bullmq"
import { connection } from "./redis" // 复用同一个 ioredis 实例
const emailQueue = new Queue("email", { connection })
new Worker(
"email",
async (job) => {
await sendEmail(job.data.to, job.data.subject)
},
{ connection },
)
await emailQueue.add("welcome", { to: "a@b.com", subject: "欢迎" })
2.3 连接复用
Queue 与 Worker 都需要 Redis 连接,但 Worker 会阻塞等待任务,因此必须使用独立的连接实例(BullMQ 内部会自行创建),不要把同一个 ioredis 实例同时交给 Queue 与 Worker。
2.4 任务状态机
一个任务在 waiting → active → completed 之间流转,失败进入 failed,等待重试则回到 delayed,超过重试次数则永久停留在 failed。理解状态机是排查「任务卡住」类问题的前提。
一句话总结:Queue 生产、Worker 消费、Job 承载状态,QueueEvents 订阅事件——Worker 必须用独立连接,理解状态机是排查卡任务的基础。
3. 任务类型与 payload 泛型
3.1 用泛型声明队列契约
type EmailJob =
| { name: "welcome"; data: { to: string; subject: string } }
| { name: "reset"; data: { to: string; token: string } }
type EmailQueue = Queue<EmailJob["data"], void, EmailJob["name"]>
BullMQ 的 Queue<DataType, ReturnType, NameType> 三个泛型分别对应任务数据、返回值与任务名。把三者都约束住,add 时数据与名字不匹配会直接编译报错。
3.2 类型安全的生产者
const queue = new Queue<EmailJob["data"], void, EmailJob["name"]>("email", { connection })
await queue.add("welcome", { to: "a@b.com", subject: "欢迎" }) // OK
await queue.add("welcome", { to: "a@b.com" }) // 编译错误:缺 subject
3.3 类型安全的消费者
new Worker<EmailJob["data"], void, EmailJob["name"]>(
"email",
async (job) => {
switch (job.name) {
case "welcome": return sendWelcome(job.data) // data 收窄为 welcome 的载荷
case "reset": return sendReset(job.data)
}
},
{ connection },
)
3.4 一个队列还是一个任务一种队列
任务名共享同一队列便于统一限流与监控;任务名过多会导致 Worker 内出现巨大 switch。折中做法是按资源或速率限制分组:同一队列内任务共享并发与限流设置,跨组再拆队列。
一句话总结:用泛型把「任务名 ↔ 载荷」绑定在队列类型上——生产者与消费者都获得编译期约束,队列按资源与限流需求分组而非按任务名爆炸。
4. 重试与退避策略
4.1 声明式重试
await queue.add("call-partner", payload, {
attempts: 5,
backoff: { type: "exponential", delay: 1000 }, // 1s,2s,4s,8s...
removeOnComplete: 1000, // 只保留最近 1000 条已完成
removeOnFail: false, // 失败的保留,便于排查
})
4.2 自定义退避
new Worker("call-partner", handler, {
connection,
settings: {
backoffStrategy: (attemptsMade) =>
Math.min(60_000, 2 ** attemptsMade * 1000) + Math.random() * 1000,
},
})
加抖动可避免大量任务在同一时刻重试形成尖峰。
4.3 区分可重试与不可重试错误
不是所有失败都值得重试:网络超时、限流(429)适合重试;参数错误、权限不足(400/403)重试只会浪费配额。用自定义错误类型标记,遇到不可重试错误时直接 throw new UnrecoverableError(msg) 让 BullMQ 跳过剩余重试。
4.4 重试的坑
attempts 默认是 1(不重试),忘记设置会让失败任务直接进 failed;退避无抖动造成惊群;失败任务无限保留会撑爆 Redis;重试次数过多会让一条坏任务反复冲击下游。
一句话总结:重试用
attempts加指数退避,并加抖动——区分可重试与不可重试错误,用UnrecoverableError提前终止,同时给失败任务设置保留上限。
5. 幂等与任务去重
5.1 至少一次投递
队列的投递语义是至少一次:Worker 执行成功但还没来得及标记完成就崩溃,任务会被重新投递。因此每个任务处理器都必须是幂等的。
5.2 用 jobId 做去重
await queue.add("welcome", payload, {
jobId: `welcome:${userId}`, // 相同 jobId 不会重复入队
})
BullMQ 用 jobId 去重:相同 jobId 的任务在队列中只存在一条,这能有效防止「用户连点两次按钮触发两封邮件」。
5.3 幂等键与业务状态
async function handler(job: Job<PayPayload>) {
const done = await db.processed.findUnique({ where: { key: job.id! } })
if (done) return done.result // 已处理,直接返回
const result = await doPayment(job.data)
await db.processed.create({ data: { key: job.id!, result } })
return result
}
5.4 去重与幂等的边界
jobId 去重只在任务仍在队列中时生效,任务完成后即可重新入队;业务幂等键需要自己持久化。两者互补:jobId 防重复入队,业务幂等键防重复执行。
一句话总结:队列是至少一次投递,处理器必须幂等——
jobId防重复入队,业务侧持久化幂等键防重复执行,两者缺一不可。
6. 定时与延时任务
6.1 延时任务
// 15 分钟后检查订单是否支付
await queue.add("check-paid", { orderId }, { delay: 15 * 60 * 1000 })
延时任务进入 delayed 状态,到点后转为 waiting。它适合「超时取消」「稍后重试」这类场景,比自建 setTimeout 可靠得多——进程重启后延时任务依然存在。
6.2 周期任务
await queue.add("daily-report", {}, {
repeat: { pattern: "0 6 * * *" }, // 每天 6 点
jobId: "daily-report", // 避免重复注册
})
周期任务(repeatable)会在 Redis 中注册调度计划,多个实例同时注册相同 jobId 的重复任务时只保留一份,因此多副本部署下不会重复触发。
6.3 周期任务的坑
修改 pattern 不会自动删除旧的调度计划,需要先 removeRepeatable;周期任务不应依赖精确触发时间(可能有秒级抖动);周期任务若执行时间超过间隔会堆积,要配合并发与限流。
6.4 分布式定时
不要把定时逻辑写进应用进程的 setInterval——多副本会重复执行,进程重启会丢失。所有周期性工作都应注册为队列的重复任务,由 Redis 统一调度。
一句话总结:延时与周期任务都应交给队列而非进程内定时器——
delay保证进程重启不丢,repeat配合固定jobId保证多副本不重复触发。
7. 失败告警与死信队列
7.1 监听失败事件
const events = new QueueEvents("email", { connection })
events.on("failed", async ({ jobId, failedReason }) => {
logger.error({ jobId, failedReason }, "job failed")
await metrics.incr("queue_failed_total", { queue: "email" })
})
7.2 死信队列
重试耗尽后任务永久停在 failed。生产上通常把这类任务搬运到死信队列(DLQ)单独存储,避免与正常队列混淆,也便于人工介入:
events.on("failed", async ({ jobId }) => {
const job = await queue.getJob(jobId)
if (job && job.attemptsMade >= (job.opts.attempts ?? 1)) {
await dlq.add(job.name, { ...job.data, _failedReason: job.failedReason })
}
})
7.3 告警策略
失败率超阈值告警、waiting 队列持续增长告警、active 任务长时间不结束告警(可能卡死)、DLQ 有新条目告警。告警要带 jobId 与 failedReason,否则只能人工翻日志。
7.4 重放与补偿
DLQ 中的任务要能人工重放:修完 bug 后把任务重新 add 回原队列。因此死信里必须保留完整的原始载荷与上下文,而不是只存一个错误字符串。
一句话总结:用
failed事件接告警、把重试耗尽的任务搬进死信队列——死信必须保留完整载荷以便重放,告警要覆盖失败率、队列堆积与卡死任务。
8. Worker 并发与限流
8.1 并发度
new Worker("email", handler, {
connection,
concurrency: 20, // 同时最多处理 20 个任务
})
concurrency 决定同一 Worker 并行处理的任务数。对 I/O 密集任务可调高,对 CPU 密集任务调高反而因事件循环争抢而变慢,此时应改用多进程或多线程。
8.2 限流
new Worker("call-partner", handler, {
connection,
limiter: { max: 10, duration: 1000 }, // 每秒最多 10 个
})
限流用于保护下游:第三方 API 有 QPS 配额、数据库连接有限、外部系统脆弱。限流是在 Worker 侧全局生效的,跨实例共享。
8.3 优先级
await queue.add("urgent", payload, { priority: 1 }) // 数字越小越优先
await queue.add("bulk", payload, { priority: 100 })
优先级让紧急任务插队,但要注意:大量低优先级任务堆积时,高优先级任务仍能被及时处理,而反过来高优先级持续涌入会饿死低优先级任务。
8.4 多进程与多机
单机多进程可绕过单进程 CPU 限制;多机部署时所有 Worker 共享同一 Redis 队列,天然负载均衡。要避免的是「同一任务被两个 Worker 同时取走」——BullMQ 用 Lua 脚本保证取任务的原子性,无需业务侧加锁。
一句话总结:并发度按 I/O 或 CPU 密集度设定,限流保护下游——优先级防饿死要谨慎,多机扩展由 Redis 原子取任务天然支持。
9. 可观测与任务追踪
9.1 关键指标
| 指标 | 含义 | 告警信号 |
|---|---|---|
| 队列深度 | 待处理任务数 | 持续增长 |
| 处理速率 | 每秒完成任务数 | 骤降 |
| 失败率 | 失败/总数 | 超阈值 |
| 任务耗时 | 单任务执行时长 | P99 上升 |
| 重试次数 | 平均重试次数 | 上升 |
9.2 用 traceId 串联
把上游请求的 traceId 放进任务数据,Worker 执行时取出并作为 Span 的父上下文,就能在追踪系统里看到「请求 → 入队 → 出队执行」的完整链路:
await queue.add("send", { ...payload, traceId: currentTraceId() })
// Worker 内
const ctx = propagation.extract(context.active(), { traceparent: job.data.traceId })
9.3 任务上下文日志
每条日志都带上 jobId、queue、attempt,这样在日志系统里按 jobId 聚合就能看到一条任务的全部尝试记录,比按时间翻日志高效得多。
9.4 面板与巡检
用 Bull Board 或自建看板展示各队列的 waiting/active/failed 数量;日常巡检关注「failed 数量是否清零」「delayed 是否按时出队」「active 是否有长时间不结束的任务」。
一句话总结:队列深度、失败率、任务耗时是最该盯的三个指标——把
traceId塞进载荷串联全链路,日志带上jobId与attempt便于按任务聚合。
10. 优雅关闭与生产实践
10.1 优雅关闭
const worker = new Worker("email", handler, { connection })
async function shutdown() {
await worker.close() // 停止取新任务,等待进行中的任务完成
await queue.close()
await connection.quit()
}
process.on("SIGTERM", () => void shutdown())
不优雅关闭会导致进行中的任务被中断,进而触发重试,产生重复执行。
10.2 部署与滚动更新
滚动更新时新旧 Worker 会短暂共存,此时队列中的任务由双方共同消费,属于正常现象;要避免的是旧版本 Worker 消费到新版本才支持的任务名,因此任务协议变更要先兼容再上线。
10.3 踩坑清单
忘记 attempts 导致不重试;退避无抖动造成惊群;jobId 复用导致任务被静默丢弃;失败任务无限保留撑爆 Redis;Worker 与 Queue 共用连接导致阻塞;未优雅关闭造成重复执行;removeOnComplete: false 让完成记录无限堆积。
10.4 容量与容量规划
按「峰值任务数 ÷ 单 Worker 吞吐」估算 Worker 数量,并留 2 倍余量应对重试与突发;Redis 内存要按「队列深度 × 单任务大小」估算,给失败任务设置明确的保留策略。
一句话总结:生产化 = 幂等处理器 + 有界重试 + 死信告警 + 限流保护 + 优雅关闭——任务协议要先兼容后上线,Redis 保留策略与容量都要显式规划。
延伸阅读
- Node 后端的服务化实践
- 异步控制与并发治理
- 载荷的运行时校验
- 追踪与指标接入
- 微服务与异步任务拆分
- TypeScript 专题 — TypeScript 专题
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。