本节目标:把「失败」当成一等公民来设计。你会掌握 BullMQ 的重试与退避配置、如何用
UnrecoverableError及时止损、幂等键的两种落地方案,以及死信队列从搬移到重放的完整闭环。读完本节,你应该能回答一个具体问题:这条 job 失败了三次,接下来它去哪,谁负责知道。
9.2 重试、幂等与死信
上一节我们把 job 投进了队列。投递成功给人一种「事情办完了」的错觉,但真实系统里,worker 会宕机、下游会超时、数据库会死锁。本节讨论的全是「失败之后」。
9.2.1 先给失败分类
重试策略的正确与否,取决于失败的性质。把下面这张表贴在团队 wiki 上比记住 API 更有用:
| 失败类型 | 例子 | 该不该重试 |
|---|---|---|
| 瞬时故障 | 网络抖动、连接池耗尽、下游 503 | 应该,且要退避 |
| 限流 | 第三方 API 返回 429 | 应该,且要尊重 Retry-After |
| 死锁/冲突 | 数据库序列化失败 | 应该,立即重试往往就成功 |
| 数据错误 | payload 字段缺失、格式非法 | 不应该,重试一万次也一样 |
| 业务拒绝 | 邮箱不存在、余额不足 | 不应该,属于正常结果 |
最糟的做法是无差别重试。一条 payload 格式错的 job 重试 20 次,只是在放大噪音、浪费 worker 配额,还掩盖了真正的 bug。
9.2.2 attempts 与 backoff
BullMQ 把重试配置放在 job 的 opts 里:
await emailQueue.add('welcome', payload, {
attempts: 5,
backoff: { type: 'exponential', delay: 1000 },
removeOnComplete: { count: 1000 },
removeOnFail: false, // 失败的要留着,死信处理要用
});
backoff 支持三种形式:
| type | 第 n 次重试的等待时间 | 适用场景 |
|---|---|---|
fixed | 恒为 delay | 本地资源竞争、数据库死锁 |
exponential | delay * 2^(n-1) | 绝大多数外部依赖 |
| 自定义 | 由 backoffStrategies 决定 | 需要加 jitter 或读 Retry-After |
按上面的配置,一次 job 的时间线是:
| 事件 | 时间点 |
|---|---|
| 第一次执行 | T+0s |
| 第 2 次(失败后等 1s) | T+1s |
| 第 3 次(等 2s) | T+3s |
| 第 4 次(等 4s) | T+7s |
| 第 5 次(等 8s) | T+15s |
| 5 次用尽 → 进入 failed | T+15s |
指数退避的问题在于「同步性」:同一批 job 会一起失败、一起重试,形成周期性脉冲。生产环境应当加抖动:
const withJitter = (delay: number) => delay * (0.5 + Math.random());
await queue.add(name, data, {
attempts: 5,
backoff: { type: 'exponential', delay: 1000 },
// 自定义策略需要注册在 Worker 上,见下
});
// worker 侧注册自定义退避策略
new Worker('tasks', processor, {
connection,
settings: {
backoffStrategy: (attemptsMade, type, err) => {
if (type !== 'jitter') return -1; // 不处理则交回内置策略
return withJitter(1000 * 2 ** (attemptsMade - 1));
},
},
});
9.2.3 用 UnrecoverableError 及时止损
对第 9.2.1 节里「不该重试」的失败,BullMQ 提供了专门的错误类型:
import { UnrecoverableError, type Job } from 'bullmq';
import { ZodError } from 'zod';
async function processEmail(job: Job<EmailPayload, EmailResult>) {
const parsed = EmailPayloadSchema.safeParse(job.data);
if (!parsed.success) {
// 数据错了,重试无意义 —— 直接判死
throw new UnrecoverableError(`invalid payload: ${parsed.error.message}`);
}
const res = await mailer.send(parsed.data);
if (res.status === 429) {
throw new Error('rate limited'); // 可重试
}
return { messageId: res.id };
}
抛 UnrecoverableError 后,BullMQ 不会再消耗 attempts,job 立刻进入 failed。这是最容易被忽视的一个 API——很多团队的重试次数都浪费在了永远不会成功的数据错误上。
9.2.4 至少一次:重试为什么要求幂等
这是本节最重要的认知:BullMQ 的投递语义是 at-least-once(至少一次),不是 exactly-once。
原因在于执行与确认是两步,而它们之间有一个无法消除的窗口:
worker 取出 job → 执行 processor → 写回 completed
↑
进程在这里崩溃
如果崩溃发生在「副作用已经产生」和「写回 completed」之间,job 会被判定为失败并重新投递。于是 processor 被第二次执行——发了两封邮件、扣了两次库存、重复创建了订单。
这不是 BullMQ 的实现缺陷,而是分布式系统的普遍结论:跨越网络的「恰好一次」无法仅靠重试实现,必须由业务侧的幂等来兜底。这也是为什么生产者在投递时通常要带一个业务侧的唯一标识:
await enqueue('order:charge', { orderId, idempotencyKey: `charge:${orderId}` });
9.2.5 幂等键的两种落地方式
方案一:Redis SETNX 去重。 适合短窗口去重(比如 24 小时内不重复),成本低:
import { redis } from './redis';
export async function once<T>(key: string, ttlSec: number, fn: () => Promise<T>): Promise<T | null> {
const acquired = await redis.set(key, '1', 'EX', ttlSec, 'NX');
if (!acquired) {
return null; // 已经执行过,直接跳过
}
try {
return await fn();
} catch (err) {
await redis.del(key); // 执行失败要释放,否则重试会被误判为重复
throw err;
}
}
new Worker('tasks', async (job) => {
const { orderId } = job.data;
return once(`charge:${orderId}`, 86400, async () => {
return payments.charge(orderId);
});
});
注意那个 catch 里的 del:幂等标记必须在业务成功后才保留。如果先占坑再执行、执行失败却不释放,那么一次失败就会让后续所有重试都被「去重」掉——任务静默丢失,比重复执行更危险。
方案二:数据库唯一索引。 适合需要永久幂等或要留痕的场景:
await prisma.idempotencyRecord.create({
data: { key: `charge:${orderId}`, jobId: job.id! },
});
| 维度 | Redis SETNX | 数据库唯一索引 |
|---|---|---|
| 生效窗口 | 受 TTL 限制 | 永久 |
| 性能 | 高 | 受写入影响 |
| 事务一致性 | 与业务事务难对齐 | 可与业务同事务 |
| 可追溯 | 差(无记录) | 好(有表可查) |
| 适用 | 通知、报表等可容忍重复 | 支付、库存等资金相关 |
对资金相关的操作,务必用方案二并与业务写入放在同一个事务里;否则「幂等记录写了但业务没写」或反之,都会造成状态不一致。更系统的讨论可参考站内专题 分布式幂等与可靠性设计 与 幂等模式在分布式系统中的落地 。
9.2.6 死信队列:失败之后去哪
BullMQ 的 failed 集合不是死信队列。它只是一个「失败 job 的索引」,默认还会被 removeOnFail 清理,而且没人会主动去看它。真正的死信队列应当是另一个队列,承载那些「重试已用尽、需要人处理」的 job:
import { Queue, QueueEvents, Worker } from 'bullmq';
const tasks = new Queue('tasks', { connection });
const dlq = new Queue('tasks:dlq', { connection });
// 监听失败事件,把耗尽重试的 job 搬进死信队列
const events = new QueueEvents('tasks', { connection });
events.on('failed', async ({ jobId, failedReason }) => {
const job = await tasks.getJob(jobId);
if (!job) return;
const maxAttempts = job.opts.attempts ?? 1;
if (job.attemptsMade < maxAttempts) return; // 还会重试,先不搬
await dlq.add(job.name, {
originalJobId: jobId,
payload: job.data,
reason: failedReason,
failedAt: new Date().toISOString(),
attemptsMade: job.attemptsMade,
});
});
| 维度 | BullMQ 的 failed | 独立 DLQ |
|---|---|---|
| 本质 | job 索引(还在原队列) | 独立的队列,可单独消费 |
| 生命周期 | 受 removeOnFail 影响 | 由你决定,通常长期保留 |
| 谁能看到 | 查日志的人 | 告警系统、值班的人 |
| 能否重放 | 需手动 retry() | 天然支持(重新投递) |
| 定位 | 调试线索 | 处理闭环 |
搬移时保留原始 payload 与失败原因,是让死信「可重放」的关键——只搬一个 id,等原始 job 被清理后就彻底丢了。
9.2.7 死信的处理闭环
死信队列如果没人消费,和没有一样。一个可用的闭环至少包含三件事:告警、人工处理、批量重放。
// 死信积压告警
setInterval(async () => {
const waiting = await dlq.getWaitingCount();
if (waiting > 100) {
logger.error({ waiting }, 'DLQ backlog too high');
await alerting.notify('#oncall', `死信队列积压 ${waiting} 条`);
}
}, 60_000);
// 重放脚本:修好数据后把死信重新投回原队列
import { Queue } from 'bullmq';
const dlq = new Queue('tasks:dlq', { connection });
const tasks = new Queue('tasks', { connection });
const jobs = await dlq.getJobs(['waiting', 'delayed'], 0, 100);
for (const job of jobs) {
const { payload } = job.data as { payload: unknown };
await tasks.add(job.name, payload, { attempts: 5 });
await job.remove();
}
console.log(`replayed ${jobs.length} jobs`);
$ node scripts/replay-dlq.ts
replayed 37 jobs
重放的前提是 processor 已经修好——否则只是把同样的失败再演一遍,而且会污染 attemptsMade 的统计。上线重放脚本前,先在预发环境用同一批死信数据跑一遍。
9.2.8 四个常见坑
一、重试次数与退避不匹配。 attempts: 20 配 fixed 100ms 等于对下游发起 20 次连击,故障时会把对方打得更死。重试的目的是等待恢复,不是加大压力。
二、幂等键用了随机值。 idempotencyKey: crypto.randomUUID() 每次都不一样,等于没有幂等。键必须由业务数据确定性生成,比如 charge:${orderId}。
三、副作用不可回滚却仍盲目重试。 调用外部支付、发短信这类操作,一旦发出就无法撤销。这类 job 应当把「发起」与「确认」拆成两步,或者依赖对方提供的幂等键(多数支付网关都支持 Idempotency-Key 头)。
四、日志里没有 jobId。 出问题时排查的第一句话通常是「这条 job 到底跑了几次」。在 processor 开头就把 job.id、job.name、job.attemptsMade 打进结构化日志,成本极低、收益极高,具体写法见 《TypeScript编程实战》3.3 结构化日志与脱敏
。
9.2.9 与本书其它章节的衔接
上一节 《TypeScript编程实战》9.1 BullMQ 队列模型与 payload 泛型
建立的契约类型,正是本节 zod 校验与幂等键的来源。想了解如何用 Result 表达「可重试 / 不可重试」的返回值而不是靠抛错,可读 《TypeScript编程实战》3.1 Result/Either 与类型化错误
。worker 进程收到 SIGTERM 时如何停止接收新 job、把在途 job 跑完,见 《TypeScript编程实战》5.3 优雅关闭与健康检查
;死信积压与重试率的指标接入见 《TypeScript编程实战》17.2 指标与告警
。
延伸阅读:消息队列的高级特性:死信与重试 、Webhook 的重试与幂等 、工作流中的重试与幂等 、Go 超时与退避重试的实践 。
小结
本节的核心结论只有一句:队列给你的是至少一次,幂等由你自己负责。
围绕这句话展开的是三件事。重试策略上,先用 UnrecoverableError 把不该重试的失败摘出去,再对可重试的失败用指数退避加抖动,避免同步脉冲。幂等设计上,幂等键必须由业务数据确定性生成,Redis SETNX 适合短窗口去重、数据库唯一索引适合资金相关场景,且失败时必须释放标记,否则重试会被误判为重复。死信处理上,BullMQ 的 failed 集合只是索引不是队列,真正的 DLQ 是另一个队列,搬移时要带上原始 payload 与失败原因,并且必须有告警与重放脚本才算闭环。
重试与死信解决的是「一次任务失败怎么办」。但还有一类任务压根不该由人触发——它们按时间表周期性运行,而且并发数必须被严格控制,否则定时任务会变成定时雪崩。这就是下一节 《TypeScript编程实战》9.3 定时任务与并发限流 的内容。
阅读导航:上一节:9.1 BullMQ 队列模型与 payload 泛型 · 下一节:9.3 定时任务与并发限流 。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。