首页 / 文章 / 使用 BullMQ 和 Redis 构建可靠的后台任务系统

使用 BullMQ 和 Redis 构建可靠的后台任务系统

学习如何使用 BullMQ 和 Redis 设计具有弹性的 Node.js 后台任务处理流程,涵盖重试、并发、幂等性以及监控等内容。

2933 词

发送确认邮件、生成报告、处理支付——在回复用户之前,有许多后台工作并不需要立即完成。本指南将介绍如何利用 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处理任务管理,工作进程则执行实际的加工操作。重试机制用于应对临时性故障,并发设置有助于控制处理速度,监控功能则能及时发现异常情况。而精心设计的应用层架构可确保在必要时任务能够安全地多次运行。

    这里的核心要点是:

    并非所有问题都必须在请求-响应周期内解决。

    有时最恰当的回应仅仅是:

    “我已经接收了这项任务,其余的交给我们处理。”

    相关阅读