Skip to content

消息队列与异步任务

用户提交一个操作,后端要发邮件、生成报表、同步外部系统——这些事如果放在请求里同步做,用户就要等它们全部完成才能看到响应,其中任何一个慢或失败都会拖垮整个请求。异步任务的核心思路是:把"不需要立即完成的工作"移出请求路径,交给后台可靠地处理,让请求快速返回。

什么该进队列

判断标准很简单:用户是否需要立即看到这个操作的结果

放在请求里(同步)放进队列(异步)
创建订单的核心写入、扣库存发送确认邮件、推送通知
返回用户必须立即看到的数据生成报表、导出文件
权限校验、支付扣款同步数据到第三方、日志归档

一个常见误区是把所有耗时操作都异步化。如果用户在等待这个结果(比如支付),强行异步反而让体验更差——你要么让用户轮询,要么推送结果,复杂度都上来了。原则是:用户在等的就同步做完,用户能等的才异步化

队列的基本结构

一个异步任务系统由三部分组成:

text
生产者:在请求中创建任务,写入队列,立即返回

队列:暂存待处理的任务(内存/数据库/专业 broker)

消费者:后台进程从队列取出任务,执行,确认完成
js
// 请求里:创建任务就返回,不等执行
async function registerHandler(req) {
  const user = await userService.create(req.body)
  await emailQueue.add('welcome', { userId: user.id })  // 入队,不等待发送
  return { ok: true }  // 用户立即得到响应
}

// 后台消费者:处理邮件发送
emailQueue.process('welcome', async (job) => {
  await sendWelcomeEmail(job.data.userId)
})

队列的存在让生产者和消费者解耦:请求不需要知道邮件怎么发、发多久;消费者也不需要知道是谁触发的。任何一边变慢或重启,另一边不受影响。

可靠性的三个保证

异步任务真正的难点不是"跑起来",而是"在各种故障下仍能完成":

1. 任务不丢

如果入队后、执行前服务崩了,任务不能消失。用持久化队列(任务存数据库或专业 broker),生产者写入成功才算入队,消费者取走后不立即删除,而是标记为"处理中",处理完成才删除。这样崩溃后能从"处理中"恢复。

2. 失败重试

任务会失败——网络超时、外部服务暂时不可用、数据冲突。需要一个重试机制:

text
执行失败 → 延迟重试(指数退避)→ 达到上限 → 进入死信队列/人工处理

重试要有上限和退避,否则一个永远失败的任务会无限重试,占满消费者。无法自动恢复的任务,进入死信队列等人工介入,而不是悄悄丢弃。

3. 幂等执行

任务可能被重复执行——消费者处理完了但确认前崩了,重启后会再取一次。所以任务的执行逻辑必须是幂等的:

text
发邮件任务 → 先查"这封邮件发过吗" → 没发过才发
创建订单任务 → 用幂等键判断"这个请求处理过吗" → 处理过就跳过

不带幂等的重试,是产生重复数据(重复扣款、重复发邮件)的常见原因。

常见错误:重试不幂等,导致重复扣款

这是异步任务里代价最高的坑——钱和数据出问题:

js
// 发送扣款任务,失败就重试
paymentQueue.process('charge', async (job) => {
  await chargeCard(job.data)  // 没有任何"是否已扣过"的判断
})

症状:用户只下了一单,却收到两次扣款;只注册了一次,却收到三封欢迎邮件。最难查的是——日志显示每次重试都"成功"了,看起来一切正常,钱却多扣了。

为什么错:消费者处理完任务、但在发"完成确认"之前崩了,重启后会把这个任务再取出来执行一次。如果执行逻辑不幂等,重试就等于重复执行了一次真实的副作用(扣款、发邮件、创建订单)。

正确做法:执行体必须能识别"这件事我已经做过":

js
paymentQueue.process('charge', async (job) => {
  // 用幂等键判断是否已处理
  const exists = await db.query('SELECT 1 FROM payments WHERE idempotency_key = ?', [job.data.key])
  if (exists) return  // 已处理,跳过

  await chargeCard(job.data)
  await db.query('INSERT INTO payments (idempotency_key, ...) VALUES (?, ...)', [job.data.key])
})

发邮件就先查"这封发过吗",创建订单就用幂等键去重。重试一定会发生,执行逻辑必须假设自己会被重复调用

优先级与并发

不是所有任务都同等重要。系统通知邮件可以慢慢发,密码重置邮件应该优先;报表可以排队,支付回调要尽快处理。用优先级队列让重要任务插队。

并发数也要控制。消费者不是越多越好——下游服务(数据库、外部 API)有承受上限,消费者太多反而压垮依赖:

text
邮件服务限制 10 并发 → 消费者并发设 8,留余量
报表生成很耗 CPU → 并发设 2,避免抢占主服务资源

可观测性

异步任务看不见摸不着,没有监控就是黑盒。每个任务应该能回答:

  • 现在队列里积压了多少?(积压增长说明消费速度跟不上生产)
  • 平均处理时长多少?(突然变慢说明依赖有问题)
  • 失败率多少?(持续失败说明代码或环境有问题)
  • 有没有卡在"处理中"很久的任务?(可能是消费者挂了没确认)
text
关键指标:
  队列深度(待处理数量)
  处理延迟(从入队到完成)
  失败率与重试次数
  死信队列积压

发现积压时的第一反应不应该是"加消费者",而是先搞清楚为什么积压——是流量突增,还是某个任务卡住、消费者挂了。盲目加消费者可能只是把问题推迟。

务实的设计清单

  • 这个任务用户是否需要立即看到结果?需要就同步,不需要才异步。
  • 队列是否持久化,任务在崩溃后能否恢复?
  • 失败是否有重试、上限和死信处理?
  • 任务执行是否幂等,重复消费会不会产生脏数据?
  • 是否有积压、延迟、失败的监控和告警?

异步任务是用"最终一致"换取请求的快速响应。它把工作从"现在必须完成"变成"一定会完成",但这个"一定"需要持久化、重试和幂等共同保证。把这些基础设施做扎实,异步化才真正可靠,而不是把同步的问题变成更难查的异步问题。

为复用而记录,为理解而整理。