ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

NestJS异步任务实战:用Bull+Redis构建可靠消息队列

NestJS异步任务实战:用Bull+Redis构建可靠消息队列 说实话做后端开发写到一定阶段一定会遇到这么个场景用户注册完要在同一请求里带着发邮件、初始化默认配置、甚至还要打个埋点一套同步流程全部跑完。测试环境问题不大一旦上了生产请求一多接口直接堆延迟Redis、数据库连接池跟着被拖垮用户那边只看到一直转圈。问题不是“代码能不能跑”而是“这些事该不该在这个请求周期里干”。这种时候我们就需要把“必须立刻完成的”和“可以稍后再做的”拆开。前者留在主流程后者丢进一个可靠的中转站由后台慢慢消化。今天要聊的就是 NestJS 里最常见的异步任务与消息队列方案——Bull Redis。这篇文章不是讲概念是会带上完整的代码、配置、部署注意事项以及我实际跑生产时踩过的坑。1. 先理清楚为什么这个场景非要引入消息队列1.1 从一次用户注册说起你写一个/user/register接口逻辑大概是这样校验参数写入用户表发送一封欢迎邮件给用户初始化一个默认空文件夹之类的资源返回“注册成功”第 3、4 步如果直接放在接口函数里同步执行每次注册请求都要等邮件服务器响应、等初始化逻辑跑完哪怕每个操作只要 200ms用户就会觉得“怎么注册这么慢”。更麻烦的是如果邮件服务临时宕机整个注册接口直接抛异常数据库里用户已经写进去了但用户看到的是“注册失败”。用消息队列之后做法就变成写入用户表之后立刻返回“注册成功”同时往队列里丢一个任务send-welcome-email。后台 Worker 收到任务再慢慢发邮件。快和稳两件事同时做到了。1.2 直接把任务放在“网络请求”里为什么不行有人会问我用 NestJS 的setTimeout或者进程内的EventEmitter把耗时逻辑延后执行不行吗小流量场景下完全行大流量场景下会出现几个现实问题进程重启任务就没了setTimeout也好、进程内事件队列也好全在内存里。代码发布、服务器重启未执行的任务直接消失。多实例部署重复执行应用水平扩容成 3 个实例之后每个实例里的定时器都在跑同一个任务可能被执行 3 次。没有重试机制邮件服务失败任务就是失败了没有“等一下再试”的机会。没办法延迟执行比如“用户下单后 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 \ --restartalways \ 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 InsightRedis 官方出的 GUI 工具。我能看到 Key 的层级结构、内存占用还能直接执行命令排查问题时非常直观。Another Redis Desktop Manager开源社区作品可以查看 Redis Streams、List、ZSet 等数据类型轻量好用个人认为偶尔调试比官方工具更顺手。2.3 NestJS 接入依赖安装与模块注册接入 Bull 需要安装两个包nestjs/bull和bull。注意这里有一个我见过不少人踩的坑Bull 的主版本迭代较快nestjs/bull对不同主版本的支持不同。目前主流的组合是nestjs/bull10bullmq或者旧一点的nestjs/bull0.xbull。写这篇文章的时候我推荐用当前覆盖面最广的稳定组合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 提供三种常用配置项配合重试attempts总尝试次数。1 表示不重试3 表示首次执行加两次重试。backoff重试退避策略。可以设固定延迟{ type: fixed, delay: 5000 }或者指数退避{ type: exponential, delay: 1000 }。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 判定失败于是触发重试但邮件服务器实际已经收到了请求、发出去了。你重试一次用户就收到两封邮件。解决方案只有一个消费端幂等。说白了就是任务处理代码要能接受“同一份数据被处理两次但业务结果等价”。实际做法分两类自然幂等比如“更新用户最后登录时间为 xxx”执行两次结果一致这种不用特殊处理。需要人工幂等键比如发邮件、发短信在任务数据里带一个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 脚本能保证这两步是一个整体。如果不想重复造轮子直接上更成熟的方案RedlockNode 生态里有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 项目里引入异步任务以上这套配置直接照着配就行把队列名替换成你的业务名注意幂等就足够支撑大部分企业级场景了。
返回列表