15 - 幂等性与可靠任务设计
任务队列不能简单理解为“业务代码绝对只执行一次”。BullMQ 的目标是尽量做到一次投递,但官方明确指出,在最坏情况下可能出现至少一次投递。重试、Worker 崩溃、任务锁丢失等情况都可能让同一业务动作再次执行。
因此,可靠系统的核心不是假设“不会重复”,而是让重复执行仍然得到正确结果。
1. 什么是幂等任务
一个任务执行一次或执行多次,系统最终状态相同,就称它具有幂等性。
天然接近幂等的操作:
把订单状态设置为 shipped
非幂等操作:
库存减 1
账户扣款 100 元
发送一封邮件
后者每执行一次都会产生新的副作用,必须增加去重或状态保护。
2. 重复执行通常怎样发生
典型时间线如下:
- Worker 调用外部支付接口,支付已经成功。
- Worker 在把任务标记为 completed 前崩溃或失去 Redis 连接。
- 任务锁过期,BullMQ 将 stalled 任务重新放回等待状态。
- 另一个 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。处理任务时:
- 开启数据库事务。
- 尝试插入操作记录。
- 如果唯一键已存在,读取之前结果并直接返回。
- 在同一事务中完成业务更新。
- 提交事务。
数据库唯一约束比“先 SELECT 再 INSERT”的应用层判断更可靠,因为两个 Worker 可能同时通过 SELECT 检查。这里的原子性应由数据库事务和约束保证,而不是 Redis 分布式锁的时间假设。
6. 将复杂工作拆成分步骤任务
假设“生成月报”需要:
- 查询并汇总数据。
- 生成 PDF。
- 上传文件。
- 发送通知。
不要把四步塞进一个巨大处理函数。可以使用 Flow 把它们拆为多个可观察、可重试的任务,或者保存明确的业务状态机。
每一步都要独立幂等:
- 汇总:相同报告月份覆盖同一份草稿。
- 生成文件:写入临时文件后原子重命名,或使用固定对象 key。
- 上传:使用报告 ID 作为对象 key,重复上传覆盖相同版本。
- 通知:通过通知记录的唯一键防止重复发送。
拆分任务不是为了形式上的“微服务化”,而是为了缩小一次失败的影响范围。
7. 记录阶段状态的编码技巧
type ReportStage = 'created' | 'generated' | 'uploaded' | 'notified';
interface ReportRecord {
id: string;
stage: ReportStage;
storageKey: string | null;
}
Worker 开始某一步前先读取持久化状态:如果目标阶段已经完成就返回已有结果;否则执行操作并保存新状态。状态必须存储在数据库或可靠外部系统中,不能只放在 Worker 内存里。
需要特别避免这种写法:
let hasSentEmail = false;
进程重启后内存标志会丢失,也无法在多个 Worker 间共享。
8. 重试与错误分类
- 网络抖动、上游 503:通常抛出普通
Error,配合attempts和 backoff 重试。 - 数据格式永久错误、业务明确拒绝:使用
UnrecoverableError,避免无意义重试。 - 部分成功:先查询外部系统或本地步骤状态,再决定继续、补偿或返回已有结果。
不要通过捕获异常后直接返回成功来隐藏失败,否则 BullMQ 会把任务标记为 completed,监控和重试都会失效。
9. 最佳实践清单
- 默认假设任务可能重复执行。
- 使用稳定的业务幂等键,并在真正产生副作用的系统中校验它。
- 用唯一约束和事务处理并发竞争。
- 让 job data 保存标识和必要参数,不保存不可验证的临时内存状态。
- 复杂任务拆分为简单步骤,每一步独立幂等。
- 记录外部系统返回的资源 ID,重试时先查询。
- 区分可恢复错误和永久错误。
- 对“副作用成功但 Worker 崩溃”的时间窗口专门编写测试。