失败、重试与停滞任务

后台任务一定会失败。网络抖动、第三方接口超时、数据库连接中断适合重试;参数错误、资源不存在等确定性错误通常不应重试。可靠队列的重点不是“永不失败”,而是让失败可观察、可恢复,并避免重复执行造成副作用。

BullMQ 何时认为任务失败

任务会在以下情况进入 failed 集合:

处理函数必须抛出 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 退避

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}`);
}

建议先区分错误类型:

stalled 到底是什么意思

任务进入 active 后,BullMQ 会为它加锁,Worker 会定期续锁。如果 Node.js 事件循环长时间被 CPU 密集型代码占满,Worker 无法及时续锁,任务就可能被判定为 stalled

停滞不等于业务代码主动抛错。BullMQ 会把停滞任务移回等待状态,使其可能被另一个 Worker 再次执行;超过最大停滞次数后才进入 failed。因此,同一任务可能已经产生部分副作用,又被执行一次。

常见原因:

实用处理方式:

幂等是重试安全的基础

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 并发时仍可能重复。优先使用数据库唯一约束、第三方接口的幂等键,或原子状态迁移。

最佳实践清单

本章小结

attemptsbackoff 解决暂时性故障,UnrecoverableError 用来停止无意义的重试;stalled 表示任务锁未能续期,任务可能被再次处理。可靠性的最后一道防线不是更多重试,而是幂等、监控和清晰的错误分类。

官方文档