失败、重试与停滞任务
后台任务一定会失败。网络抖动、第三方接口超时、数据库连接中断适合重试;参数错误、资源不存在等确定性错误通常不应重试。可靠队列的重点不是“永不失败”,而是让失败可观察、可恢复,并避免重复执行造成副作用。
BullMQ 何时认为任务失败
任务会在以下情况进入 failed 集合:
- Worker 的处理函数抛出异常。
- 任务反复停滞,超过 Worker 的
maxStalledCount。
处理函数必须抛出 Error 对象,不要抛字符串:
throw new Error("生成报表失败");
使用 attempts 自动重试
attempts 表示允许执行的总次数,包含第一次执行。attempts: 3 即首次失败后最多再重试两次。
import { Queue } from "bullmq";
interface ReportJobData {
reportId: string;
}
const reportQueue = new Queue<ReportJobData>("report", {
connection: { host: "127.0.0.1", port: 6379 },
defaultJobOptions: {
attempts: 4,
backoff: {
type: "exponential",
delay: 1_000,
jitter: 0.2,
},
removeOnComplete: 500,
removeOnFail: 2_000,
},
});
await reportQueue.add("generate", { reportId: "report-42" });
如果没有配置 backoff,失败任务会尽快重新排队。生产环境通常应加入退避,避免故障期间不断冲击下游服务。
fixed 与 exponential 退避
fixed:每次等待固定时长,适合短暂、恢复时间相对稳定的故障。exponential:等待时间按指数增长,计算方式约为2^(attemptsMade - 1) * delay,更适合外部 API 或数据库故障。jitter:在等待时间中加入随机性,取值为0到1,可避免大量任务同时重试形成“惊群”。
await reportQueue.add(
"generate",
{ reportId: "report-43" },
{
attempts: 3,
backoff: { type: "fixed", delay: 5_000, jitter: 0.5 },
},
);
重试任务重新进入等待状态时仍会遵守自身优先级。
监听失败事件
import { Job, Worker } from "bullmq";
interface ReportJobData {
reportId: string;
}
const worker = new Worker<ReportJobData>(
"report",
async (job: Job<ReportJobData>): Promise<void> => {
await generateReport(job.data.reportId);
},
{ connection: { host: "127.0.0.1", port: 6379 } },
);
worker.on("failed", (job, error: Error) => {
console.error("任务失败", {
jobId: job?.id,
attemptsMade: job?.attemptsMade,
message: error.message,
});
});
async function generateReport(reportId: string): Promise<void> {
console.log(`生成报表:${reportId}`);
}
failed 事件适合记录日志、指标和告警,不适合依赖它完成必须成功的业务写入,因为事件监听器本身也可能失败。
不可恢复错误
确定继续重试也不会成功时,抛出 UnrecoverableError。它会覆盖任务的 attempts 配置,直接停止后续重试。
import { UnrecoverableError, Worker } from "bullmq";
interface EmailJobData {
address: string;
}
const worker = new Worker<EmailJobData>(
"email",
async (job): Promise<void> => {
if (!job.data.address.includes("@")) {
throw new UnrecoverableError("邮箱地址格式无效");
}
await sendEmail(job.data.address);
},
{ connection: { host: "127.0.0.1", port: 6379 } },
);
async function sendEmail(address: string): Promise<void> {
console.log(`发送邮件到 ${address}`);
}
建议先区分错误类型:
- 超时、连接重置、HTTP 503:通常可重试。
- 参数校验失败、明确的 HTTP 400/404:通常不可恢复。
- HTTP 429:属于限流,应按服务端给出的等待时间处理,参见第 10 章。
stalled 到底是什么意思
任务进入 active 后,BullMQ 会为它加锁,Worker 会定期续锁。如果 Node.js 事件循环长时间被 CPU 密集型代码占满,Worker 无法及时续锁,任务就可能被判定为 stalled。
停滞不等于业务代码主动抛错。BullMQ 会把停滞任务移回等待状态,使其可能被另一个 Worker 再次执行;超过最大停滞次数后才进入 failed。因此,同一任务可能已经产生部分副作用,又被执行一次。
常见原因:
- 同步执行大规模 JSON 解析、压缩或加密。
- 无限循环或耗时很长的同步算法。
- Worker 进程崩溃、被强制终止或机器失联。
实用处理方式:
- 将 CPU 密集型处理拆小,或放到独立进程 / Worker Threads。
- 不要随意缩短
stalledInterval;官方说明通常无需修改它。 - 为停滞和失败次数建立监控,而不是只看“队列有没有消费”。
- 优雅关闭 Worker,让当前任务有机会处理完毕。
幂等是重试安全的基础
BullMQ 的任务处理应按“至少一次”思维设计:任务可能因为重试、停滞恢复或进程故障而重复执行。
interface ChargeJobData {
orderId: string;
amountInCents: number;
}
async function chargeOnce(data: ChargeJobData): Promise<void> {
// 实际项目应在数据库中为 orderId 建立唯一约束,
// 并在同一事务中检查和记录扣款结果。
const alreadyCharged = await hasChargeRecord(data.orderId);
if (alreadyCharged) return;
await chargePayment(data);
await saveChargeRecord(data.orderId);
}
不要只依赖“先查询、再写入”的内存判断;多 Worker 并发时仍可能重复。优先使用数据库唯一约束、第三方接口的幂等键,或原子状态迁移。
最佳实践清单
- 明确哪些错误可重试,哪些错误应直接失败。
- 默认使用有限次数的指数退避并加入
jitter。 - 只抛
Error或其子类。 - 保留适量失败任务,便于排查,但设置
removeOnFail防止无限增长。 - 日志至少包含队列名、任务 ID、任务名、尝试次数和错误信息。
- 对发邮件、扣款、创建订单等有副作用的操作保证幂等。
- 不把 CPU 密集型同步工作直接塞进普通 Worker 的事件循环。
本章小结
attempts 和 backoff 解决暂时性故障,UnrecoverableError 用来停止无意义的重试;stalled 表示任务锁未能续期,任务可能被再次处理。可靠性的最后一道防线不是更多重试,而是幂等、监控和清晰的错误分类。