使用 BullMQ 和 Redis 构建可靠的后台任务系统
学习如何使用 BullMQ 和 Redis 设计具有弹性的 Node.js 后台任务处理流程,涵盖重试、并发、幂等性以及监控等内容。
发送确认邮件、生成报告、处理支付——在回复用户之前,有许多后台工作并不需要立即完成。本指南将介绍如何利用 BullMQ 与 Redis 结合来构建可靠的后台任务系统。
想象这样一个后台系统:几乎所有任务都在 HTTP 请求周期内直接完成。需要发送邮件?就在那里处理。需要生成 PDF?也是同样方式。需要处理一些后台数据?也在请求中直接完成。
这种方法起初运作良好,但后来就会出问题。
API 开始变慢,请求开始超时。如果某个外部服务出现故障,整个请求也会随之失败。
正是在这种情况下,后台任务才显得十分重要。
不必强制 API 在回复之前完成所有步骤,而是可以将任务放入队列中,由专门的处理程序单独处理。
在 Node.js 生态系统中,一个不错的选择是BullMQ,它以 Redis 作为存储后端。以下是其工作原理。
1. 什么是后台任务?
后台任务指的是那些无需作为 HTTP 请求的一部分同步执行的操作。
以典型的注册流程为例。当有人创建账户时,API 可能需要执行以下操作:
- 创建用户记录
- 发送欢迎邮件
- 生成欢迎 PDF
- 发送通知
- 更新其他下游系统
你可能会尝试在回复之前直接执行所有这些操作:
Client
↓
API
↓
Create User
↓
Send Email
↓
Generate PDF
↓
Send Notification
↓
Response
但这样会迫使用户等待所有步骤都完成。
更好的方法应该是这样:
Client
↓
API
↓
Create User
↓
Add Job to Queue
↓
Response
另外:
Queue
↓
Worker
↓
Send Email
↓
Done
因为API不再需要在回复之前完成所有任务,所以它的响应速度会快很多。
2. 为什么我们需要队列?
假设发送邮件需要1秒,生成PDF需要2秒,调用另一个API也需要1秒。那么你的接口在发送响应之前可能会停滞数秒——这对用户来说体验很差。
更糟糕的是,如果邮件服务不可用会怎样?即使创建用户的过程实际上已经成功,请求仍可能失败。这是两个无关功能之间不必要的依赖关系。
队列可以打破这种耦合:
┌──────────────┐
│ Node API │
└──────┬───────┘
↓
Add Job
↓
┌──────────────┐
│ Redis │
│ Queue │
└──────┬───────┘
↓
┌──────────────┐
│ Worker │
└──────┬───────┘
↓
Email / PDF / API / etc.
通过这样的分离,API和后台任务各自承担明确的职责。
3. 什么是BullMQ?
BullMQ 是一个用于 Node.js 的队列库,它依赖 Redis 来存储和协调任务。其架构从高层次上看如下所示:
Producer
↓
Queue
↓
Worker
↓
Job Processing
生产者负责创建任务,队列则用于存储这些任务,而工作进程则负责实际处理它们。
例如:
await emailQueue.add("welcome-email", {
userId: user.id,
email: user.email
});
本质上,该 API 的作用是:
“有一些工作需要完成。”
它本身并不需要执行这些工作。
4. 创建队列
一个最简单的 BullMQ 队列配置如下所示:
import { Queue } from "bullmq";
const connection = {
host: "localhost",
port: 6379
};const emailQueue = new Queue("email", {
connection
});
之后,你可以将任务推送到该队列中:
await emailQueue.add("welcome-email", {
userId: "123",
email: "user@example.com"
});
Redis 在后台负责存储所有与队列相关的状态。从概念上讲,可以这样理解:
email queue
Job 1
Job 2
Job 3
Job 4
Job 5
然后,工作进程会获取并处理这些任务。
5. 创建工作进程
工作进程才是实际执行任务的组件:
import { Worker } from "bullmq";
const worker = new Worker(
"email",
async (job) => {
console.log("Processing:", job.name); await sendWelcomeEmail(
job.data.email
);
},
{
connection
}
);
综合起来,整个流程现在如下所示:
API
↓
emailQueue.add()
↓
Redis
↓
Worker
↓
sendWelcomeEmail()
API无需等待邮件发送完成——而这正是后台任务带来的核心优势。
6. 当任务失败时会发生什么?
这正是队列相比普通服务调用展现出真正优势的地方。
想象这样的场景:
API
↓
Email Service
↓
ERROR
当直接调用API时,你必须立即决定如何处理失败情况。
而队列则提供了另一种选择:可以简单地重新尝试执行该任务。
以下是一个示例:
await emailQueue.add(
"welcome-email",
{
email: "user@example.com"
},
{
attempts: 3
}
);
通过这种配置,任务在放弃之前可以尝试多次。
从流程图上看,其结构如下:
Attempt 1
↓
Failed
↓
Attempt 2
↓
Failed
↓
Attempt 3
↓
Success
在处理不可靠的第三方服务时,这种模式极具价值。
不过,重试次数不应没有限制或随意进行。
应避免出现任务无止境地重复重试且没有终止条件的情况。
7. 带延迟的重试
假设某个外部服务暂时出现故障。
需要避免的情况如下:
FAIL
RETRY IMMEDIATELY
FAIL
RETRY IMMEDIATELY
FAIL
RETRY IMMEDIATELY
立即反复重试正在出问题的服务实际上可能会使状况恶化。
解决办法是在每次尝试之间加入延迟。
例如:
await emailQueue.add(
"welcome-email",
{
email: "user@example.com"
},
{
attempts: 5,
backoff: {
type: "exponential",
delay: 5000
}
}
);
从概念上来看,其结构如下:
Attempt 1 → Fail
↓
5 sec
↓
Attempt 2 → Fail
↓
10 sec
↓
Attempt 3 → Fail
↓
20 sec
↓
Attempt 4 → Success
具体的时间间隔取决于重试及延迟策略的配置方式。
但核心理念始终不变:
让临时故障有时间恢复后再尝试。
8. 延迟执行的任务
并非所有任务都需要在创建后立即运行。
例如:
在注册24小时后发送提醒。
BullMQ允许你安排任务在稍后执行。
await emailQueue.add(
"reminder",
{
userId: "123"
},
{
delay: 24 * 60 * 60 * 1000
}
);
从概念上讲:
Create Job
↓
Wait 24 hours
↓
Worker processes job
这种模式常见于以下场景:
- 提醒邮件
- 定时通知
- 试用期限结束
- 付款提醒
- 后续沟通信息
9. 多个工作进程
现在想象一个每分钟接收数千个任务的系统。
单个工作进程可能无法应对如此高的负载。
你可以通过同时运行多个工作进程来实现扩展:
Redis Queue
↓
┌──────────┼──────────┐
↓ ↓ ↓
Worker 1 Worker 2 Worker 3
↓ ↓ ↓
Jobs Jobs Jobs
每个进程都会独立地从队列中获取任务。
例如:
1000 email jobs
你可能会看到类似这样的情况:
Worker 1 → Job 1, 4, 7...
Worker 2 → Job 2, 5, 8...
Worker 3 → Job 3, 6, 9...
增加更多工作进程是提升处理效率的一种方法。
但需注意:
单纯增加工作进程并不一定就能带来好处。
你的数据库、邮件服务提供商、CPU、内存以及所有下游服务都有各自的容量限制。
10. 并发性
除了运行多个工作进程外,BullMQ还允许你配置单个工作进程同时处理的任务数量。
例如:
const worker = new Worker(
"email",
async (job) => {
await sendEmail(job.data.email);
},
{
connection,
concurrency: 5
}
);
这样可以让一个工作进程并行处理多个任务。
从概念上讲:
Worker
├── Job 1
├── Job 2
├── Job 3
├── Job 4
└── Job 5
更高的并发性可以提高处理效率。
但不要不加思考就将并发度直接提升到100。
如果每个任务都会访问数据库,过高的并发度很容易使其超载。
并发设置应根据实际工作负载能力来调整。
11. 速率限制
有时瓶颈根本不在你自己的系统中,而在于你所依赖的第三方服务。
假设你的邮件服务提供商限制了每秒的请求次数。
如果你突然有:
10,000 jobs
你就不应该一次性发送所有请求。
队列可以控制任务处理的速率。
最终的架构如下:
10,000 Jobs
↓
Queue
↓
Rate Limit
↓
Worker
↓
External API
这种方式比向服务提供商发送数千个同步请求要安全得多。
12. 任务幂等性很重要
下一个概念是后台任务处理中最关键的思路之一。
以支付处理任务为例:
Process Payment
工作进程负责执行该任务。
支付成功完成。
但就在工作进程将其标记为已完成之前,进程突然崩溃。
队列按照设计功能重新尝试执行该任务。
如果没有防护措施,就可能导致向客户再次收费。
这是个非常严重且代价高昂的问题。
为避免这种情况,应在可行范围内让任务具备幂等性。
实际上,这意味着重复执行同一任务不应产生任何意外的重复副作用。
一种常见的方法是以唯一的支付参考编号作为关键标识:
payment:order_123
然后,在执行任何操作之前先进行检查:
Has this payment already been completed?
↓
Yes → Don't charge again
↓
No → Process payment
BullMQ 本身没有为此内置的机制。
必须通过应用程序代码来实现幂等性。
13. 失败的任务需要相应的策略
并非所有的失败情况都相同,也并非所有失败都需要重试。
以下是一些示例:
Invalid email
Invalid user ID
Missing database record
Invalid payment information
再次运行这些任务五次也无法解决问题。
将失败情况分为两类会有所帮助:
临时性故障
这类故障包括:
- 网络超时
- 某个依赖项暂时不可用
- 数据库连接中断
对于这类问题,稍后再次尝试确实是合理的。
永久性故障
这类故障包括:
- 无效的输入数据
- 已不存在的引用资源
对于这类情况,重新尝试毫无意义——任务应直接进入某种故障处理流程。
一个设计良好的队列系统并非简单地遵循以下通用规则:
Retry everything
而是采用更为周全的处理流程:
Understand why it failed
↓
Temporary?
/ \
YES NO
↓ ↓
Retry Handle failure
14. 无法处理的任务/失败任务处理
无论你多么小心,总有一些任务会以无法通过重试解决的方式失败。你需要能够查看这些任务,避免它们直接消失。
例如,可能会出现如下情况:
Failed Jobs
──────────────
Job 101 → Email invalid
Job 102 → Payment failed
Job 103 → API timeout
一旦能够看到这些故障,你就有以下处理选项:
- 记录故障以便日后审查
- 通知团队成员
- 让某人能够手动重新尝试
- 修正导致故障的错误数据
具体的实现方式取决于系统的需求。最重要的是遵循一个原则:
出错的任务绝不能毫无痕迹地消失。
15. 队列与定时任务
这两者很容易被混淆,但实际上它们解决的是不同的问题。
定时任务的任务是:
“在特定时间执行此任务。”
队列的任务是:
“处理这个任务单元。”
在实际应用中,这两种工具往往能很好地配合使用。例如:
Cron
↓
Find users whose trial expires today
↓
Create jobs
↓
Queue
↓
Workers
↓
Send emails
这样可以将调度逻辑与处理逻辑分开。相比让单个定时进程试图独自完成所有工作,这种设计通常更为清晰。
16. 队列事件与监控
一旦在生产环境中运行该系统,就需要能够了解队列内部的实际状况。
值得跟踪的指标包括:
- 等待处理的任务
- 当前正在处理的任务
- 成功完成的任务
- 失败的任务
- 处理所需的时间
- 重试的次数
- 队列的总大小
想象一下控制面板突然显示出类似这样的内容:
Waiting Jobs
Normal: 50
Current: 25,000
这种突然的波动是一个警示信号,可能意味着:
- 处理任务的工作进程已停止运行
- 外部 API 的响应速度变慢
- 数据库负载过重
- 流量突然激增
- 最近的部署带来了错误
如果不监控队列,这些问题会悄悄累积,直到用户发现异常为止。
17. 不要所有任务都放入队列
拥有 BullMQ 并不意味着每个操作都必须作为后台任务处理。
比如以下情况:
GET /profile
在这里,用户需要立即获取他们的个人资料数据。将其放入后台队列处理毫无意义——只会增加不必要的延迟。
只有在以下情况下才适合使用队列:
- 任务需要较长时间才能完成
- 任务可以异步执行
- 任务可能需要重试
- 任务占用大量资源
- 任务依赖于不可完全依赖的外部服务
- 结果无需包含在即时响应中
一个有用的问题是:
在返回HTTP响应之前,用户真的需要这个结果吗?
如果不需要,可以考虑将相关处理放到后台任务中执行。
18. 生产环境风格的架构
将所有组件整合在一起后,典型的架构如下所示:
Client
↓
Node.js API
↓
┌──────┴──────┐
↓ ↓
PostgreSQL Redis
↓
Queue
↓
┌──────────┼──────────┐
↓ ↓ ↓
Worker 1 Worker 2 Worker 3
↓ ↓ ↓
Email PDF Notifications
API层负责处理那些需要立即处理的操作。PostgreSQL(或您选择的数据库)用于存储持久性的业务数据。Redis则负责队列基础设施以及那些适合其处理的短期任务。工作进程则处理所有可以异步完成的任务。
以这种方式划分职责,能让整个系统的扩展变得容易得多。
19. 应避免的错误
错误1:在HTTP请求中处理所有操作
这会导致API既慢又脆弱。
错误2:无限制地重试
有些故障无论尝试多少次都无法自行解决。
错误3:忽略幂等性
如果某个任务被执行两次,可能会引发你未曾预料的重复副作用。
错误4:允许无限并发
如果没有限制,就有可能让任务所依赖的系统不堪重负。
错误5:忽略监控
若队列持续无限制地增长,就会成为即将出现的运营问题。
错误6:将Redis用作系统记录存储
队列状态与核心业务数据用途不同,不应混为一谈。
错误7:让所有操作都异步进行
有些操作确实需要在完成之后才能返回响应。
20. 更好的思维模型
在理解队列之前,人们的本能反应往往是这样思考请求处理的:
Request
↓
Do everything
↓
Response
而一个更有用的模型应该是这样的:
Request
↓
Do what must happen immediately
↓
Queue what can happen later
↓
Response
接着是:
Queue
↓
Worker
↓
Process
↓
Retry if appropriate
↓
Complete / Fail
这种将必须立即执行的任务与可以稍后处理的任务分开的做法,正是这一切背后的核心思想。
最终结论
BullMQ 的价值并非仅仅因为它是一款广泛使用的 Node.js 库,而是因为后台任务处理能够满足真正的架构需求。
如果某项工作具有以下特点:
- 处理速度慢
- 可以重试
- 不会阻塞响应生成
- 依赖于外部服务
- 需要大量资源
那么它很可能就不应该放在 HTTP 请求处理流程中。
队列为这些任务提供了执行场所。Redis负责提供底层基础设施,BullMQ处理任务管理,工作进程则执行实际的加工操作。重试机制用于应对临时性故障,并发设置有助于控制处理速度,监控功能则能及时发现异常情况。而精心设计的应用层架构可确保在必要时任务能够安全地多次运行。
这里的核心要点是:
并非所有问题都必须在请求-响应周期内解决。
有时最恰当的回应仅仅是:
“我已经接收了这项任务,其余的交给我们处理。”
相关阅读
- 设计实时聊天后端:房间管理、数据持久化与扩展性 — 了解如何使用 Socket.IO、PostgreSQL 和 Redis 构建实时聊天后端,涵盖房间管理、消息持久化顺序、在线状态检测以及多服务器扩展等内容。
- Redis 缓存基础:设计模式、常见陷阱与面试问题 — 学习 Redis 缓存在 Node.js 应用中的工作原理,包括缓存旁路策略、TTL 设置、防止挤兑机制、数据驱逐策略以及常见的面试问题。