15 - 幂等性与可靠任务设计

任务队列不能简单理解为“业务代码绝对只执行一次”。BullMQ 的目标是尽量做到一次投递,但官方明确指出,在最坏情况下可能出现至少一次投递。重试、Worker 崩溃、任务锁丢失等情况都可能让同一业务动作再次执行。

因此,可靠系统的核心不是假设“不会重复”,而是让重复执行仍然得到正确结果。

1. 什么是幂等任务

一个任务执行一次或执行多次,系统最终状态相同,就称它具有幂等性。

天然接近幂等的操作:

把订单状态设置为 shipped

非幂等操作:

库存减 1
账户扣款 100 元
发送一封邮件

后者每执行一次都会产生新的副作用,必须增加去重或状态保护。

2. 重复执行通常怎样发生

典型时间线如下:

  1. Worker 调用外部支付接口,支付已经成功。
  2. Worker 在把任务标记为 completed 前崩溃或失去 Redis 连接。
  3. 任务锁过期,BullMQ 将 stalled 任务重新放回等待状态。
  4. 另一个 Worker 再次处理任务。

BullMQ 无法自动撤销第 1 步,也无法仅凭 Redis 状态知道外部系统是否已经成功。因此业务侧必须提供幂等键或可恢复状态。

3. 官方建议:任务保持简单、原子

BullMQ 的 idempotent jobs pattern 建议任务尽量简单和原子。一个任务同时更新数据库、上传文件、调用支付、发送邮件,会产生很多“完成了一半”的状态,也很难安全回滚。

更适合重试的任务通常具备以下特点:

4. 用业务幂等键保护副作用

jobId 可以避免同一个 ID 的任务被重复加入队列,但它不能单独解决所有重复执行问题:已存在任务的清理策略、不同队列的 ID 空间以及 Worker 已开始执行后的崩溃都需要考虑。

真正涉及扣款、发券等操作时,应把幂等键传给负责副作用的系统:

import { Job, Worker } from 'bullmq';

interface ChargeJobData {
  amountInCents: number;
  orderId: string;
}

interface ChargeResult {
  paymentId: string;
}

interface PaymentGateway {
  charge(input: {
    amountInCents: number;
    idempotencyKey: string;
  }): Promise<ChargeResult>;
}

declare const paymentGateway: PaymentGateway;

const worker = new Worker<ChargeJobData, ChargeResult>(
  'payments',
  async (job: Job<ChargeJobData>): Promise<ChargeResult> => {
    return paymentGateway.charge({
      amountInCents: job.data.amountInCents,
      idempotencyKey: `charge:${job.data.orderId}`,
    });
  },
  { connection: { host: '127.0.0.1', port: 6379 } },
);

幂等键应来自稳定的业务标识,例如订单号,而不是每次重试都会变化的随机 UUID。

5. 数据库中的“只执行一次”记录

常见模式是在数据库建立唯一约束,例如 operation_key UNIQUE。处理任务时:

  1. 开启数据库事务。
  2. 尝试插入操作记录。
  3. 如果唯一键已存在,读取之前结果并直接返回。
  4. 在同一事务中完成业务更新。
  5. 提交事务。

数据库唯一约束比“先 SELECT 再 INSERT”的应用层判断更可靠,因为两个 Worker 可能同时通过 SELECT 检查。这里的原子性应由数据库事务和约束保证,而不是 Redis 分布式锁的时间假设。

6. 将复杂工作拆成分步骤任务

假设“生成月报”需要:

  1. 查询并汇总数据。
  2. 生成 PDF。
  3. 上传文件。
  4. 发送通知。

不要把四步塞进一个巨大处理函数。可以使用 Flow 把它们拆为多个可观察、可重试的任务,或者保存明确的业务状态机。

每一步都要独立幂等:

拆分任务不是为了形式上的“微服务化”,而是为了缩小一次失败的影响范围。

7. 记录阶段状态的编码技巧

type ReportStage = 'created' | 'generated' | 'uploaded' | 'notified';

interface ReportRecord {
  id: string;
  stage: ReportStage;
  storageKey: string | null;
}

Worker 开始某一步前先读取持久化状态:如果目标阶段已经完成就返回已有结果;否则执行操作并保存新状态。状态必须存储在数据库或可靠外部系统中,不能只放在 Worker 内存里。

需要特别避免这种写法:

let hasSentEmail = false;

进程重启后内存标志会丢失,也无法在多个 Worker 间共享。

8. 重试与错误分类

不要通过捕获异常后直接返回成功来隐藏失败,否则 BullMQ 会把任务标记为 completed,监控和重试都会失效。

9. 最佳实践清单

官方资料