☰
NestJS异步任务实战:用Bull+Redis构建可靠消息队列
2026/10/7 16:58:24 网站建设 项目流程

说实话,做后端开发写到一定阶段,一定会遇到这么个场景:用户注册完,要在同一请求里带着发邮件、初始化默认配置、甚至还要打个埋点,一套同步流程全部跑完。测试环境问题不大,一旦上了生产,请求一多,接口直接堆延迟,Redis、数据库连接池跟着被拖垮,用户那边只看到一直转圈。问题不是“代码能不能跑”,而是“这些事该不该在这个请求周期里干”。

这种时候,我们就需要把“必须立刻完成的”和“可以稍后再做的”拆开。前者留在主流程,后者丢进一个可靠的中转站,由后台慢慢消化。今天要聊的,就是 NestJS 里最常见的异步任务与消息队列方案——Bull + Redis。这篇文章不是讲概念,是会带上完整的代码、配置、部署注意事项,以及我实际跑生产时踩过的坑。


1. 先理清楚:为什么这个场景非要引入消息队列

1.1 从一次用户注册说起

你写一个/user/register接口,逻辑大概是这样:

  1. 校验参数
  2. 写入用户表
  3. 发送一封欢迎邮件
  4. 给用户初始化一个默认空文件夹之类的资源
  5. 返回“注册成功”

第 3、4 步如果直接放在接口函数里同步执行,每次注册请求都要等邮件服务器响应、等初始化逻辑跑完,哪怕每个操作只要 200ms,用户就会觉得“怎么注册这么慢”。更麻烦的是,如果邮件服务临时宕机,整个注册接口直接抛异常,数据库里用户已经写进去了,但用户看到的是“注册失败”。

用消息队列之后,做法就变成:写入用户表之后,立刻返回“注册成功”,同时往队列里丢一个任务:send-welcome-email。后台 Worker 收到任务再慢慢发邮件。快和稳,两件事同时做到了。

1.2 直接把任务放在“网络请求”里为什么不行

有人会问:我用 NestJS 的setTimeout或者进程内的EventEmitter把耗时逻辑延后执行,不行吗?

小流量场景下完全行,大流量场景下会出现几个现实问题:

  1. 进程重启,任务就没了:setTimeout也好、进程内事件队列也好,全在内存里。代码发布、服务器重启,未执行的任务直接消失。
  2. 多实例部署重复执行:应用水平扩容成 3 个实例之后,每个实例里的定时器都在跑,同一个任务可能被执行 3 次。
  3. 没有重试机制:邮件服务失败,任务就是失败了,没有“等一下再试”的机会。
  4. 没办法延迟执行:比如“用户下单后 30 分钟未支付自动关闭订单”,这种延时任务用setTimeout写,一重启就全部失效。

所以我们需要一个外部的、有持久化能力的队列中间件。任务写进去之后,不依赖某个进程的存活,并且支持延迟、重试、去抖、并发控制这些语义。这时候 Redis + Bull 就是个非常合适的组合。

1.3 Bull 而不是别的队列:选型逻辑

消息队列领域有 RabbitMQ、Kafka、RocketMQ 这些重型选手,那为什么在 NestJS 里往往首选 Bull?

关键原因是语义匹配度。Bull 是专为 Node.js 设计的、基于 Redis 的轻量级任务队列库,它实现的不是通用消息流,而是“Job”模型——一个任务,有生命周期、有进度、有重试策略、有超时控制。对 Web 应用常见的异步任务(发邮件、生成 PDF、推送通知、定时汇总)来说,这个模型比 Kafka 的分区消费模型直观得多。

而且 Bull 对 NestJS 有原生支持。官方提供了@nestjs/bull封装,程序员只需要声明一个队列、写一个处理器,模块会自动完成 Redis 连接、消费注册、事件监听这些繁琐工作,不需要手写 Redis 的 BLPOP 循环。我形容它为兔子,Kafka 是鲸鱼,如果你的业务还没到“每天上亿条日志流”,Bull 的运维成本和理解成本都低得多。等到真有海量数据流需求,再升级不迟。

另外,Bull 底层依赖的 Redis 同时也是你项目里的缓存中间件,不需要额外加一套服务,省了一笔运维成本。这正是它在中小型团队和中等规模项目中比 RabbitMQ 更普及的根本原因。


2. 环境准备:Redis 是基础设施,先把它跑起来

2.1 安装 Redis:三种环境一次说清

Boot 是队列存储层,所以第一步是准备一个可用的 Redis 实例。这里把开发环境最常遇到的三个场景一次说清楚:

MacOS 用户

用 Homebrew 是最省事的:

brew tap redis/redis brew install redis brew services start redis

安装完成后用redis-cli ping,返回PONG就说明服务正常。

Windows 用户

Redis 官方并没有原生 Windows 版本,但微软之前维护过分支,社区也一直有移植版本。目前最稳的路线是用 WSL 或者 Docker。如果你不想折腾,也可以到 Redis 官方 GitHub 仓库的 releases 页面找 Windows 版压缩包,解压后直接运行redis-server.exe。有个很实用的小技巧:把解压目录加入系统 PATH,这样你在任意目录都能用redis-cli。

Docker 用户

选 Docker 的好处是环境能统一,特别是多机协作时保证版本差异不坑人:

docker run -d \ --name redis \ -p 6379:6379 \ --restart=always \ redis:7.2-alpine

提示:生产环境不要这样直接裸跑。至少加一个密码:redis-server --requirepass yourpassword。如果你是在云上跑 Redis,且安全组没限制来源 IP,那相当于把数据裸奔在公网,这是最容易被忽略的安全问题。

2.2 连接自检与可视化工具

装好之后,需要用客户端确认连接情况。很多人死在这一步:明明 Redis 起了,NestJS 却连接不上,最后发现是密码错了,或者是在配置文件里多打了空格。

常见连接检测命令:

redis-cli -h 127.0.0.1 -p 6379 redis-cli -a yourpassword ping

如果带了密码,redis-cli -a命令会提示 warning,这个不用紧张,只是提醒你密码出现在命令行里可能存在日志泄露风险。本地开发无所谓,生产环境建议用REDISCLI_AUTH环境变量传递密码。

想看队列的数据结构,推荐两个工具:

  • Redis Insight:Redis 官方出的 GUI 工具。我能看到 Key 的层级结构、内存占用,还能直接执行命令,排查问题时非常直观。
  • Another Redis Desktop Manager:开源社区作品,可以查看 Redis Streams、List、ZSet 等数据类型,轻量好用,个人认为偶尔调试比官方工具更顺手。

2.3 NestJS 接入:依赖安装与模块注册

接入 Bull 需要安装两个包:@nestjs/bull和bull。注意,这里有一个我见过不少人踩的坑:Bull 的主版本迭代较快,@nestjs/bull对不同主版本的支持不同。目前主流的组合是@nestjs/bull@10+bullmq,或者旧一点的@nestjs/bull@0.x+bull。写这篇文章的时候,我推荐用当前覆盖面最广的稳定组合:bull和@nestjs/bull。

npm install @nestjs/bull bull npm install -D @types/bull

另外需要安装 Redis 客户端。NestJS 底层默认依赖ioredis,所以还要:

npm install ioredis

依赖装好后,在AppModule里注册 Bull 的核心模块:

import { Module } from '@nestjs/common'; import { BullModule } from '@nestjs/bull'; @Module({ imports: [ BullModule.forRoot({ redis: { host: '127.0.0.1', port: 6379, password: 'yourpassword', }, }), ], }) export class AppModule {}

这里forRoot配置的是全局默认连接。如果未来要连多个 Redis 实例,可以给不同队列分别指定连接名称,那是进阶玩法,刚起步不需要。

此时,我强烈建议先写一个最简单的队列测试,确认整条链路能跑通,再往业务里加代码。别一口气把几十个文件写完了再启动,到时候报错了你都不知道错在哪一层。


3. 核心配置:生产者与消费者的一次握手

3.1 注册一个具体的业务队列

在模块里用BullModule.registerQueue()声明这个模块要用到的队列。比如我需要一个邮件队列:

import { Module } from '@nestjs/common'; import { BullModule } from '@nestjs/bull'; import { EmailService } from './email.service'; import { EmailConsumer } from './email.processor'; @Module({ imports: [ BullModule.registerQueue({ name: 'email', }), ], providers: [EmailService, EmailConsumer], }) export class EmailModule {}

这个name就是队列的名字。Bull 会在 Redis 里创建以bull:email:为前缀的一系列 Key。不需要提前“建表”,Redis 的 Key 是动态的,任务进来自然就有了。

3.2 生产者:注入队列并添加任务

生产者就是把任务塞进队列的代码。在 NestJS 中,可以直接通过@InjectQueue()装饰器注入队列实例:

import { Injectable } from '@nestjs/common'; import { InjectQueue } from '@nestjs/bull'; import { Queue } from 'bull'; @Injectable() export class EmailService { constructor( @InjectQueue('email') private readonly emailQueue: Queue, ) {} async sendWelcomeEmail(userId: number, email: string) { await this.emailQueue.add( 'welcome', { userId, email }, { attempts: 3, backoff: { type: 'exponential', delay: 2000 }, removeOnComplete: true, }, ); } }

这里add()的第一个参数是任务类型welcome,第二个参数是任务数据。attempts: 3表示失败最多重试 3 次,backoff表示指数退避——第一次失败等 2 秒,第二次失败等 4 秒,第三次失败等 8 秒,给它呼吸空间。removeOnComplete: true是任务完成之后自动从 Redis 清除,避免堆积无用数据。

3.3 消费者:处理器与会话

消费者是真正干活的角色。在 NestJS 里用@Processor()装饰器定义处理器类,用@Process()装饰器定义具体任务的处理方法:

import { Processor, Process } from '@nestjs/bull'; import { Job } from 'bull'; @Processor('email') export class EmailConsumer { @Process('welcome') async handleWelcome(job: Job<{ userId: number; email: string }>) { console.log(`开始给用户 ${job.data.userId} 发送欢迎邮件`); // 模拟发送邮件 await sendMail(job.data.email, '欢迎注册'); console.log('邮件发送完成'); } }

消费者类注册为 provider 之后,NestJS 会把它自动挂载到对应队列上。当生产者往email队列投递welcome类型任务时,这个handleWelcome方法就会被调用。

这个模型的使用者可以怎么理解呢?生产者完全不关心消费者是谁、在哪里,消费者也不关心是谁投递的任务。它们只围绕“队列名 + 任务类型”通信,这让业务模块的解耦变得非常自然。


4. 企业级功能落地:重试、延迟、重复、并发和进度

4.1 失败重试与退避策略

实际生产里面,任务失败是常态。发邮件时 SMTP 服务器超时、调第三方 API 返回 429、生成报表时数据库连接池打满……这些并不是代码逻辑错误,而是外部环境临时抖动。如果没有重试,意味着用户会收到一条永远不响应的业务反馈。

我在配置重试上一般遵循这个原则:业务上允许延迟的任务,多配几次重试;时效性极高的任务,少配甚至不配,改为告警人工介入。

Bull 提供三种常用配置项配合重试:

  1. attempts:总尝试次数。1 表示不重试,3 表示首次执行加两次重试。
  2. backoff:重试退避策略。可以设固定延迟{ type: 'fixed', delay: 5000 },或者指数退避{ type: 'exponential', delay: 1000 }。
  3. timeout:单次执行超时。比如任务 10 秒没执行完就判定失败,交给重试逻辑。

如果重试次数耗尽仍失败,任务会进入failed状态。这时候一定要有监听:

@Processor('email') export class EmailConsumer { @OnQueueFailed() onFailed(job: Job, err: Error) { console.error(`任务 ${job.id} 最终失败,原因:${err.message}`); // 此处可以通知告警服务,或者写日志表 } }

这里我推荐一个习惯:重试耗尽不要静默丢掉,至少打个警告日志,方便事后查账。

4.2 延迟任务与定时重复任务

很多业务场景需要“过一会儿再做”。比如订单创建 15 分钟后未支付就自动关闭、活动开始前 1 小时推送提醒。在 Bull 里延迟执行很简单:

await this.orderQueue.add( 'timeout-close', { orderId }, { delay: 15 * 60 * 1000, // 15分钟后执行 attempts: 2, }, );

这个功能在底层用的是 Redis 的 ZSet。Bull 会把延迟任务按时间戳排序,到了时间再把任务移动到等待队列里。实现很直接,但非常实用。

定时重复任务则用来处理“每天凌晨清理过期数据”之类的需求:

await this.reportQueue.add( 'daily-report', {}, { repeat: { cron: '0 2 * * *', // 每天凌晨两点 }, }, );

注意:使用repeat时必须给任务一个唯一的jobId,否则重复添加会报错。最稳妥的做法是显式指定:

await this.reportQueue.add( 'daily-report', {}, { repeat: { cron: '0 2 * * *' }, jobId: 'daily-report-job', }, );

4.3 并发数与限流

默认情况下,一个 Worker 进程同时处理多个任务。如果并发度太高,会把下游服务(数据库、邮件服务、第三方 API)瞬间打爆。Bull 的并发控制很朴素——在process()方法里传并发数。

在 NestJS 的@Processor装饰器里可以这样加:

@Processor('email', { concurrency: 5 }) export class EmailConsumer { @Process('welcome') async handleWelcome(job: Job) { // 最多同时处理 5 个欢迎邮件任务 } }

这个 5 并不是拍脑袋乱写的。它应该等于你的“下游系统能承受的每秒并发 / 单个任务的预计耗时”。比如邮件服务每秒能承受 10 个请求,发一封邮件平均需要 0.2 秒,那单个进程并发数 5 是安全的。如果你起了 4 个应用实例,每个实例又配 5 并发,那就是 20 并发,这个就要重新算一下下游能不能扛住。

4.4 进度上报与 Job 事件

如果异步任务耗时较长,比如批量导出 Excel、生成视频,前端就非常需要一个“正在进行到哪一步了”的提示。Bull 支持进度上报,消费端可以这样更新进度:

@Process('export') async handleExport(job: Job) { const totalRows = 10000; for (let i = 0; i < totalRows; i += 1000) { // 批量查询数据并写入 Excel 文件 await job.progress((i / totalRows) * 100); } }

生产端可以监听任务事件拿到进度:

this.exportQueue.on('progress', (job, progress) => { // 把进度写入 Redis 缓存,前端就可以轮询获取了 await this.redis.set(`job:${job.id}:progress`, progress); });

更常用的是把进度存在 Redis 里,前端通过GET /export/status/:jobId查询。这是异步任务中最常见的一个闭环交互模式:接口提交任务 → 返回 jobId → 前端轮询进度 → 任务完成下载文件。这套流程在报表系统、数据导出、图片批量处理里几乎通用。


5. 生产环境最容易踩的坑

5.1 重复消费问题:为什么会发生,怎么解决

这是面试里高频问题,也是生产中真实会遇到的。消息队列理论上至少有三种语义:最多一次、至少一次、精确一次。Bull 默认是至少一次。

什么意思?就是任务执行过程中,如果 Worker 进程崩溃了,Bull 不知道任务是否处理到一半,等进程恢复时,任务会被重新投递。比如你的邮件服务响应超时,Bull 判定失败,于是触发重试;但邮件服务器实际已经收到了请求、发出去了。你重试一次,用户就收到两封邮件。

解决方案只有一个:消费端幂等。说白了就是任务处理代码要能接受“同一份数据被处理两次,但业务结果等价”。

实际做法分两类:

  1. 自然幂等:比如“更新用户最后登录时间为 xxx”,执行两次结果一致,这种不用特殊处理。
  2. 需要人工幂等键:比如发邮件、发短信,在任务数据里带一个messageId,执行前先查 Redis 或数据库,如果这个messageId处理过,直接跳过。
@Process('welcome') async handleWelcome(job: Job) { const { userId, email, messageId } = job.data; const key = `msg:${messageId}`; const existed = await this.redis.get(key); if (existed) { return { status: 'duplicated' }; // 防止重复发送 } await sendMail(email); await this.redis.set(key, '1', 'EX', 24 * 3600); }

这个解决方案的适用范围很广。不只是 Bull,任何说“至少一次”的消息系统,消费端幂等都是必经之路。

5.2 内存、序列化和 Redis 连接超时

Bull 在 Redis 里会保存任务数据和任务状态。如果removeOnComplete: false(默认就是 false),所有执行完的任务还留在 Redis 里,时间久了数据量会非常可观。尤其日志量大、任务频率高的时候,Redis 内存会被撑爆。

我的策略分三档:

  • 日志型任务:执行完立即删除,removeOnComplete: true, removeOnFail: false(保留失败日志)
  • 业务型任务:成功后删除,失败保留几天,手动清洗
  • 审计型场景:不删除,定期导出归档

Redis 连接超时也是高频问题。尤其在容器化环境,服务启动瞬间创建大量连接到 Redis,触达 Redis 的maxclients上限,之后的新连接就会Command timed out。排查步骤先看 Redis 日志,确认是不是连接数超过上限。预防方式是合并连接:BullModule.forRoot里不要为每个队列创建一个连接池,共享一个默认连接即可。

另外要注意:Bull 默认会把任务数据序列化为 JSON 存储在 Redis。不要在任务数据里塞Date对象、Buffer、或巨大的对象。建议只放“最小必要数据”:ID 加关键参数。消费者需要完整数据时再回源数据库查一次。有的面试题甚至是:“为什么消息体要尽量小?”答案不是节省空间,而是减少 Redis 的网络传输时间和内存压力。

5.3 分布式锁:何时需要,如何做

消息队列和分布式锁经常一起出现。前者解决异步任务分发,后者解决资源竞争。

典型场景:订单出库时要扣减库存,多个实例同时消费任务,同一件商品可能被并发扣减,数据就错了。这时候要为“每件商品的扣减动作”加一把分布式锁。

最朴素的实现方式是在 Redis 上使用SET NX EX原子操作:

import Redis from 'ioredis'; const redis = new Redis({ host: '127.0.0.1', port: 6379 }); async function acquireLock(key: string, token: string, ttl: number = 3000) { const result = await redis.set(`lock:${key}`, token, 'EX', ttl, 'NX'); return result === 'OK'; } async function releaseLock(key: string, token: string) { // 使用 Lua 脚本,保证判断和删除是原子的 const script = ` if redis.call("get", KEYS[1]) == ARGV[1] then return redis.call("del", KEYS[1]) else return 0 end `; return redis.eval(script, 1, `lock:${key}`, token); }

这里为什么要用 Lua 脚本?因为“判断 token 是否一致”和“删除 Key”是两个操作,如果不原子执行,可能你判断完还没删,锁已经过期了,另一个实例拿到了新锁,然后你把别人的锁删了。Lua 脚本能保证这两步是一个整体。

如果不想重复造轮子,直接上更成熟的方案:Redlock,Node 生态里有redlock库,简单对接一下即可。

但有一点我要提醒:分布式锁是“成本较高”的解决方案。能用乐观锁(比如数据库版本号)解决的问题,不要动不动上分布式锁,它的容错设计极其复杂,而且一旦出问题,表现非常隐蔽。

5.4 任务积压与监控

如果消费速度低于生产速度,队列里的等待任务数会越来越多。Bull 里可以这样定期看:

const counts = await this.queue.getJobCounts(); console.log(counts); // { waiting: 1000, active: 5, completed: 2000, failed: 3, delayed: 50 }

监控这个数值的意义在于提前发现问题。一般来说,waiting数量持续递增且消费实例没有扩容空间,就要排查消费者逻辑是不是出现了阻塞(比如每次请求外部 API 超时时间设置过长)。

可视化监控可以直接部署 Bull Board 的独立 UI 版本,但我更推荐先把日志和指标接进现有的监控体系。如果团队已经有 Grafana 那一套,就不要单独为了队列多起一个服务,能用日志查到问题先查日志,别让技术栈指数级膨胀。


6. 最后一点经验

写到这里,进入队列生产环境的边界也已经很清晰了。看着复杂,其实底座就一句话:用 Redis 存任务状态,用 Worker 异步消费,再配上重试、延迟和幂等,让耗时业务彻底从请求链路里剥离。

我自己在项目落地过程中,最深的体会是:刚开始上 Bull 的时候,不要一上来就想着把所有“慢操作”全部切到队列。先选一两个痛点场景(比如注册欢迎邮件、订单超时关闭)跑通闭环,观察 Redis 的内存增长和任务堆积情况,再逐渐扩大范围。原因很简单,队列化改造是有移动成本的,业务逻辑一旦拆成“生产端 + 消费端”,排查问题时的链路视野也要跟着变。先小步跑,踩熟了再大步走。

另外一个很实用的习惯分享给大家:每上线一个新的队列任务,我第一件事就是把@OnQueueFailed()里的日志打全,包括job.id、job.name、attemptsMade、完整错误堆栈。等出了问题,这些小信息能帮你节省至少一个下午的排查时间。

如果你正要在 NestJS 项目里引入异步任务,以上这套配置直接照着配就行,把队列名替换成你的业务名,注意幂等,就足够支撑大部分企业级场景了。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询