Nest 使用笔记
第二十章:RabbitMQ——把领域事件从主流程里剥出来,异步、可靠地投递
打开 learnhub 的 AmqpService,在真实代码里讲透消息队列的意义、RabbitMQ 的 Exchange/Queue/binding 模型、生产者的 fire-and-forget 模式,再补上 learnhub 没写的消费者——手动 ack、持久化、重投与死信队列(DLX)。每段代码都能在 learnhub 里指到对应文件。
- Nest
- RabbitMQ
- 消息队列
前面 TypeORM 那章里,PostService.create 在事务提交后有这两行:
// learnhub/src/modules/post/post.service.ts
this.amqp.publish('post.published', { postId: post.id, title: post.title, authorId: userId });
this.indexPost(post);
indexPost 把帖子同步进 Elasticsearch,第二十一章细讲;amqp.publish 这一行就是这一章的主角——它把「帖子发布了」这件事变成一条消息丢进 RabbitMQ,让一个独立的服务去消费(发站内信、推 feed、统计)。这一章打开 learnhub 真实在用的 AmqpService,讲清楚生产者这一侧怎么写、为什么这么写;learnhub 仓库里没有消费者(注释里写明由独立的 notification-service 消费),所以这一章的后半段会手把手写一个真实的 Nest 消费者,把旧版课程 promised 但没做的**持久化、手动 ack、重投、死信队列(DLX)**全补齐。
先搞懂:消息队列是什么 / 为什么需要 / 企业级怎么用
消息队列(MQ)是什么:一个专门用来「暂存消息」的中间件。生产者把消息丢进队列就返回,不等消费者处理完;消费者按自己的节奏从队列里取消息处理。RabbitMQ 是基于 AMQP 协议、Erlang 写的开源 MQ broker,久经生产考验。它的核心模型是 Producer → Exchange → Queue → Consumer:生产者不直接调消费者,而是把消息丢进 Exchange,按路由规则落到 Queue,消费者从 Queue 取。它不止是个队列,还有 Exchange 路由、消息确认、持久化、死信等一整套能力。
为什么需要它:很多操作同步做会拖垮服务——发站内信、生成报表、统计入库、第三方回调都是耗时活儿,HTTP 请求一直挂着等它返回,并发一上来线程池打满、DB 连接打满、整条链路雪崩。learnhub 的 PostService.create 如果改成「发帖 → 同步发完所有站内信 → 同步推 feed → 同步统计」串成一坨,任意一个下游慢一点,发帖接口就跟着慢;下游挂了发帖也跟着挂。判断口径:
- 可以异步(调用方不需要立刻拿到结果)→ 上 MQ;
- 有突发流量要削峰(瞬时并发远高于下游承受力,先丢队列再慢慢消费)→ 上 MQ;
- 上下游要解耦(生产者消费者各自演化、各自部署,互不知道对方什么时候在线)→ 上 MQ;
- 强实时同步交互(登录、下单扣库存,调用方需要立刻拿到结果决定下一步)→ 别塞 MQ 多一跳网络。
企业级怎么用(下面每一条,这一章都会在 learnhub 真实代码或消费者教学里做一遍,不是空头支票):
- 消息持久化:Queue 声明
durable: true、消息发persistent: true,broker 重启不丢——这一章第二步看到 learnhub 这么声明,第五步在消费者侧继续对齐; - 消费者手动 ack:处理完才告诉 broker「这条搞定了,可以删」,处理中途挂掉消息会被重新投递——第五步落地;
- 失败重投 + 死信队列(DLX):处理失败的消息要么 requeue 重试,要么丢进 DLX 留待人工排查——第六步落地;
- 生产者publisher confirms:broker 收到消息后回 ack 给生产者,防 broker 接收时就丢——learnhub 没做(它选了 fire-and-forget 的简单路线,第三步讲 tradeoff),这一章末尾会指出来。
这一章你会做出什么
- 打开 learnhub 的
src/modules/amqp/amqp.service.ts,逐行讲清楚生产者一侧:onModuleInit建连、assertQueue(durable: true)声明持久化队列、publish走sendToQueue(persistent: true)。 - 看
PostService.create和UserService.create怎么在事务提交后发post.published/user.registered领域事件,理解 fire-and-forget 的边界。 - 自己写一个真实的 Nest 消费者模块(learnhub 没写,这是这一章的教学 payload),消费
post.published走通完整链路。 - 在消费者里把旧版课程 promised 但没做的几件事全做掉:持久化、手动 ack、重试、死信队列(DLX)。
- 学会什么时候该用 MQ、什么时候该用 Redis pub/sub、什么时候该用 WebSocket/SSE。
前置:本机或 learnhub/docker 里能起 RabbitMQ(docker compose up rabbitmq 即可)。amqplib 已经在 learnhub 的 package.json 里了。
第一步:用 Docker 起 RabbitMQ
learnhub 的 docker/docker-compose.yml 里已经写好了 RabbitMQ 服务,带管理界面:
# learnhub/docker/docker-compose.yml
rabbitmq:
image: rabbitmq:3-management
container_name: learnhub-rabbitmq
environment:
RABBITMQ_DEFAULT_USER: guest
RABBITMQ_DEFAULT_PASS: guest
ports:
- "5672:5672" # AMQP 协议端口(代码连这个)
- "15672:15672" # 管理界面 http://localhost:15672 (guest/guest)
volumes:
- rabbitmq-data:/var/lib/rabbitmq
直接起:
cd learnhub
docker compose -f docker/docker-compose.yml up -d rabbitmq
两个端口要分清:5672 是 AMQP 协议端口,代码连这个;15672 是管理界面,浏览器打开 http://localhost:15672,guest / guest 登录,能看到 connection / channel / exchange / queue 一目了然。后面调试消息流向时,管理界面的 Queues 标签页会反复用到——看消息进没进、ack 没 ack、堆积了多少,全在这里。
rabbitmq-data 这个 volume 是 RabbitMQ 的数据目录,broker 重启后队列和持久化的消息都还在;没持久化的消息(下面会讲)重启就没了。
第二步:先看 learnhub 的 AmqpService——连接、声明队列
这是这一章的主锚点。learnhub 把 RabbitMQ 客户端封进一个 @Global() 模块,整个应用哪儿都能注入:
// learnhub/src/modules/amqp/amqp.module.ts
import { Global, Module } from '@nestjs/common';
import { AmqpService } from './amqp.service';
@Global()
@Module({
providers: [AmqpService],
exports: [AmqpService],
})
export class AmqpModule {}
@Global() 让 AmqpService 一旦在 AppModule 里 import 一次,就不需要再在每个用到它的子模块里重复 import——发事件的地方很多(post、user、评论),全标一遍 imports: [AmqpModule] 太啰嗦。这也是为什么 TypeOrmModule.forRoot 也是 @Global():基础设施性质的模块(DB、缓存、日志、MQ)通常都做成全局。
AmqpService 才是真正的肉:
// learnhub/src/modules/amqp/amqp.service.ts
import { Injectable, Logger, OnModuleInit } from '@nestjs/common';
import { connect, Channel } from 'amqplib';
const QUEUE = 'learnhub.events';
@Injectable()
export class AmqpService implements OnModuleInit {
private readonly logger = new Logger(AmqpService.name);
private channel: Channel | null = null;
async onModuleInit(): Promise<void> {
const url = process.env.RABBITMQ_URL || 'amqp://127.0.0.1:5672';
try {
const conn = await connect(url);
this.channel = await conn.createChannel();
await this.channel.assertQueue(QUEUE, { durable: true });
this.logger.log(`已连 RabbitMQ,队列 ${QUEUE}(${url})`);
} catch (e) {
this.logger.warn(`RabbitMQ 连接失败,事件将被丢弃:${(e as Error).message}`);
}
}
/** 发领域事件;channel 没就绪时丢弃并告警(业务不阻塞) */
publish(type: string, payload: unknown): void {
if (!this.channel) {
this.logger.warn(`RabbitMQ 未就绪,丢弃事件 ${type}`);
return;
}
const body = Buffer.from(JSON.stringify({ type, payload }));
this.channel.sendToQueue(QUEUE, body, { persistent: true });
}
}
几个关键点要逐个讲:
生命周期 hook OnModuleInit:Nest 提供 onModuleInit 这个钩子,在模块初始化时自动调用——比裸 Node 里手写 await connect() 然后塞进全局变量干净得多。AmqpService 实现这个接口,Nest 起来时会自动建连。连接用的 URL 从 RABBITMQ_URL 环境变量读,dev 是 amqp://127.0.0.1:5672(.env 里写的),docker compose 里写成 amqp://rabbitmq:5672(容器间用服务名访问)。
assertQueue(QUEUE, { durable: true }):assert 的意思是「确保这个队列存在,没有就建,有就幂等返回」。QUEUE 是个常量 learnhub.events——learnhub 不用复杂的 exchange 路由,所有事件都丢进这同一个队列,用消息体里的 type 字段区分事件类型,消费者侧按 type 分发。这里隐含一个 AMQP 细节:sendToQueue 不指定 exchange 时走的是默认 exchange(名字为空字符串),它会把 routing key 当成队列名直接路由——所以 learnhub 实际上没显式用 exchange,但理解 RabbitMQ 模型时要知道 exchange 这一层是存在的,第七步对比方案时也会用到。durable: true 让队列持久化:broker 重启后队列还在。注意 durable 是队列的属性,和消息的 persistent 是两回事,第五步会再对齐。
try/catch 降级:连接失败不抛错,只 warn。这是有意的:MQ 是辅助设施,不该因为 MQ 暂时不可用就让整个 Nest 应用启动不了——主流程(发帖、查帖)照样能跑,只是事件这次发不出去。这种「基础设施挂了不影响主流程」的降级思路在第十三章 Redis 那章也见过(ranking.recordView().catch(() => undefined))。
没有 @InjectConnection() 之类的装饰器:因为 amqplib 不是 Nest 原生集成的库,没有现成的 provider 可注入,learnhub 直接在 service 里 connect() 自己管连接。如果你用 @nestjs/microservices,它有 ClientProxy 那一套封装,但代价是配置更厚一层——learnhub 选了最薄的方式。
第三步:生产者——publish 为什么是 fire-and-forget
publish 方法只有四行有效代码,但藏着两个关键决策:
// learnhub/src/modules/amqp/amqp.service.ts
publish(type: string, payload: unknown): void {
if (!this.channel) {
this.logger.warn(`RabbitMQ 未就绪,丢弃事件 ${type}`);
return;
}
const body = Buffer.from(JSON.stringify({ type, payload }));
this.channel.sendToQueue(QUEUE, body, { persistent: true });
}
决策一:消息格式自封装。type + payload 包成 { type, payload } 再 JSON 序列化丢进队列。消费者拿到后 JSON.parse(msg.content) 就能按 type 分发到不同处理函数。type 实际上就是事件名(post.published / user.registered),消费者用一个 switch(type) 或 handler map 路由。这种「单队列 + 类型字段」的写法适合事件种类不多的小型系统;事件一多、消费者想按类型订阅不同队列,就该上 exchange + routing key 的正经路由(见任何 RabbitMQ 进阶资料)。
决策二:不返回 Promise、不 await。sendToQueue 本身是同步调用(amqplib 内部写到 socket 缓冲区),但不保证 broker 真的收到了。learnhub 没有开 publisher confirms(broker 收到后回 ack),所以严格说这条消息可能在网络层或 broker 端就丢了,生产者完全不知道。这是 fire-and-forget 的 tradeoff:换来了「不阻塞业务请求」的好处,代价是「broker 故障时丢消息」。
注意:fire-and-forget 在 learnhub 这种「事件丢了顶多少发一条站内信」的场景是可接受的——业务的核心真相源是 MySQL,事件只是触发副作用。但如果是「下单事件必须被消费一次」的强一致场景,绝不能这么写——要么开 publisher confirms(channel.confirmSelect() 后 publish 配合 await 写法、broker 确实收下了才往下走),要么用 outbox 模式先把消息落库再异步发。learnhub 用的是简单路线,要认清它的边界。
生产者侧的两道保险已经齐了:durable: true(队列持久化)+ persistent: true(消息持久化)。意思是:消息只要成功到了 broker,broker 重启也不丢。但「从生产者到 broker 这一跳」是没有保险的(见上一段),这是 fire-and-forget 固有的窗口期。
第四步:生产者的真实调用方——PostService.create 和 UserService.create
光看 AmqpService 不够,要看它怎么被业务调用:
// learnhub/src/modules/post/post.service.ts
async create(dto: CreatePostDto, userId: number): Promise<Post> {
const post = await this.dataSource.transaction(async (manager) => {
// ... 事务里建帖子、绑标签
});
// 事务提交成功后,再发领域事件、同步 ES
this.amqp.publish('post.published', { postId: post.id, title: post.title, authorId: userId });
this.indexPost(post);
return post;
}
// learnhub/src/modules/user/user.service.ts
const saved = await this.userRepo.save(user);
// 阶段9:发「用户注册」领域事件,由 notification-service 消费
this.amqp.publish('user.registered', { userId: saved.id, username: saved.username });
return saved;
为什么不在事务里就 publish?第八章 TypeORM 已经讲过,这里重申:事务的边界要画在对的地方。发 MQ 消息不属于「帖子创建」的原子范围——MQ 暂时不可用不该把已经成功的帖子创建也回滚。所以 learnhub 的写法是:事务先提交 → 再发事件。代价是「事务提交了、但 publish 之前进程挂了」会丢一条事件(所谓的 publish-then-commit 缺口);learnhub 接受这个缺口(fire-and-forget 路线的延伸),强一致场景要换 outbox 模式。
this.amqp 是注入的 AmqpService,构造函数里 private readonly amqp: AmqpService 一行就拿到(因为 AmqpModule 是 @Global() 的)。整个 PostService 只知道「调 publish 发个事件」,不知道 RabbitMQ 长什么样、不知道谁会消费——这就是解耦的实际样子。
思考:如果 post.published 发出去了,但消费端(notification-service)此刻挂着没消费,会发生什么?——消息会一直躺在 RabbitMQ 的 learnhub.events 队列里(因为队列 durable、消息 persistent),等消费者重新连上继续消费。这种「消费者宕机不丢消息、恢复后追上」的耐久性,就是用 MQ 比起在 PostService.create 里同步调 notificationService.send() 的最大优势——同步调用下游挂了你就直接 500 了。
第五步:消费者——learnhub 没写,我们自己写一个
AmqpService 文件头注释里写明:learnhub.events 队列由独立的 notification-service 消费。但 learnhub 主仓库里没有这个消费者。这一节我们手写一个真实的 Nest 消费者模块,把消费侧应该做的事补齐。这是教学代码,不在 learnhub 主仓库里——你可以自己开个 src/modules/notification/ 跟着写。
先建连接和监听。和生产者对称,消费侧也要建 channel、assert 同一个队列:
// 教学代码:可在自己的 notification 模块里实现
import { Injectable, OnModuleInit, Logger } from '@nestjs/common';
import { connect, Channel, ConsumeMessage } from 'amqplib';
const QUEUE = 'learnhub.events'; // 必须和生产者一致
@Injectable()
export class NotificationConsumer implements OnModuleInit {
private readonly logger = new Logger(NotificationConsumer.name);
async onModuleInit(): Promise<void> {
const url = process.env.RABBITMQ_URL || 'amqp://127.0.0.1:5672';
const conn = await connect(url);
const channel = await conn.createChannel();
await channel.assertQueue(QUEUE, { durable: true }); // 和生产者声明一致
channel.prefetch(10); // 流量控制:最多同时持有 10 条未确认消息
channel.consume(QUEUE, (msg) => this.handle(msg, channel), { noAck: false });
this.logger.log('NotificationConsumer 已启动');
}
private async handle(msg: ConsumeMessage | null, channel: Channel): Promise<void> {
if (!msg) return;
try {
const { type, payload } = JSON.parse(msg.content.toString());
if (type === 'post.published') {
await this.onPostPublished(payload); // 真正的业务:发站内信、推 feed
} else if (type === 'user.registered') {
await this.onUserRegistered(payload);
}
channel.ack(msg); // 处理成功,手动 ack,broker 才会删消息
} catch (e) {
this.logger.error(`处理失败:${(e as Error).message}`);
channel.nack(msg, false, false); // 不 requeue,进 DLX(第六步细讲)
}
}
private async onPostPublished(p: { postId: number; title: string; authorId: number }) {
this.logger.log(`收到 post.published:#${p.postId} ${p.title}`);
// 这里发邮件 / 站内信 / 推 feed,省略
}
private async onUserRegistered(p: { userId: number; username: string }) {
this.logger.log(`收到 user.registered:#${p.userId} ${p.username}`);
// 这里发欢迎邮件,省略
}
}
逐条讲消费侧必须做对的几件事:
assertQueue 必须和生产者声明一致:队列名、durable 都要对齐。AMQP 的语义是「两边任意一方声明过,另一方 assert 同一个名字、属性不同,会报错」。learnhub 生产者写的是 assertQueue('learnhub.events', { durable: true }),消费者也要一字不差。
channel.prefetch(10) 是流量控制的核心旋钮:它告诉 broker「在我 ack 之前的 10 条消息之内别再发了」。消费者最多同时持有 10 条未确认消息,处理完一条 ack 一条,broker 才会再推一条。这是「削峰」落到消费者一侧的具体机制——生产端可能瞬间涌入 1000 条,但消费端按自己的处理能力(比如 10 并发)慢慢来,DB 不会被压垮。
noAck: false + 手动 channel.ack(msg):noAck: true 是「消费者一收到 broker 就当这条处理完了」(自动确认),出错了消息直接丢;noAck: false 才是生产写法——只有你显式 ack 了,broker 才认为这条处理完了、从队列删掉。处理中途进程挂掉,broker 没收到 ack,会把这条消息重新投递给另一个消费者(或重启后的自己)。learnhub 主流程不发不收,所以没有 ack 这段代码;消费者侧必须做对。
注意:手动 ack 有一个最常见的坑——ack 必须恰好一次、且对正确的消息。重复 ack 会触发 broker 报 PRECONDITION_FAILED 关通道;ack 错消息(比如用 channel.ackAll() 或批量 ack 了更早的 tag)会顺带把还没处理完的消息也标记完成、然后丢消息。出错路径上要么 nack(msg, false, false)(不重入、进 DLX),要么 nack(msg, false, true)(重入队列尾部重试),不要静默 swallow 异常又不 ack——broker 会一直留着这条「未确认」的消息,连接断开后又再投递一次,形成无限重投。
第六步:失败处理——重投 vs 死信队列(DLX)
PostService 那侧 publish 完就完事,消息处理失败完全是消费者的事。消费者处理失败有两种合法路径:重投回原队列(requeue,给这条消息再一次机会)或丢进死信队列(DLX,留待人工排查)。这两种对应上面 catch 里的 channel.nack(msg, false, requeue),第三个布尔参数控制。
nack(msg, false, true) 是 requeue:broker 把这条消息塞回队列头部、立刻又被推出来。注意:如果失败是确定性的(比如消息格式坏了、必抛异常),requeue 会让这条消息在 broker 和消费者之间无限循环——CPU 拉满、日志刷屏,事故级别。所以 requeue 只适合瞬时性错误(DB 短暂连不上、下游 503);确定性的错误必须进 DLX。
死信队列(DLX):给主队列绑一个「死信交换机」,主队列里被 nack 且不 requeue 的消息、超过 TTL 的消息、超过最大长度的消息,都会自动转发到 DLX 上挂的队列。消费者侧不再处理,留着人工在 DLX 队列里看、修、回放。
声明带 DLX 的队列(这是这一章 promised 但旧版没做的关键能力):
// 教学代码:把第五步的 assertQueue 周围换成下面这段
const DLX = 'learnhub.events.dlx';
const DLQ = 'learnhub.events.dlq';
// 1. 声明死信交换机和死信队列(消费者侧也要建,幂等)
await channel.assertExchange(DLX, 'direct', { durable: true });
await channel.assertQueue(DLQ, { durable: true });
await channel.bindQueue(DLQ, DLX, ''); // 路由 key 留空,direct 模式下都进 DLQ
// 2. 主队列声明时挂上 x-dead-letter-exchange
await channel.assertQueue(QUEUE, {
durable: true,
arguments: {
'x-dead-letter-exchange': DLX, // 被 nack(requeue=false) 或超 TTL 的消息转给 DLX
'x-dead-letter-routing-key': '', // 转发时用的 routing key
},
});
参数对齐是关键:x-dead-letter-exchange 的值要和 assertExchange(DLX, ...) 的名字一字不差;DLX 自己可以是 direct / fanout / topic 任一种,简单场景用 direct + 空 key 就够。声明带 arguments 的队列第一次声明时就要带上——AMQP 不允许同一队列二次声明改属性,所以如果你之前已经 assertQueue(QUEUE, { durable: true }) 过(learnhub 主仓库就这么声明过),要先在管理界面里把队列删掉再重新建。
然后第五步 catch 里 nack(msg, false, false)(第三个参数 false = 不 requeue),消息就会自动路由到 DLQ。
注意:DLX 不是「自动重试 N 次后进死信」的现成机制——RabbitMQ 本身不计数重试次数。要实现「重试 3 次后进死信」,要么自己在消息头里塞 x-retry-count 每次递增、到上限就 nack(false);要么给主队列绑一个带 TTL 的「重试队列」(消息在里面等 30 秒后超时死信、再回主队列),实现「延迟重试」。生产环境常用第二种,因为加上了退避间隔。这一章不展开到那个深度,记住 DLX 是失败消息的归宿而不是自动重试机制就行。
注意:DLQ 上的消息需要监控,否则会一直堆。常见做法是给 DLQ 也挂消费者做「告警 + 落库」,或者用管理界面看 DLQ 深度配 Prometheus 告警。RabbitMQ 管理界面 Queues 标签页能看到每个队列的消息数、消费者数、ack 速率——DLQ 深度一旦上涨就说明业务在持续失败,必须人工介入。
第七步:MQ vs Redis pub/sub vs WebSocket/SSE——按需选
这一节帮你想清楚「我现在这个需求,到底该用哪个」。三种东西都能传消息,但语义完全不同:
- RabbitMQ(MQ):生产者和消费者解耦、消息持久化、消费者宕机不丢消息、有 ack/重试/DLX。适合「这件事必须被处理、但可以异步」——发站内信、订单下游同步、报表生成。典型特征:消费者不在线时消息存着等它。
- Redis pub/sub(第十三章讲过的
publish/subscribe):fire-and-forget 广播,发布那一刻有几个订阅者就投给几个,没订阅者消息直接丢,不存。适合「实时通知在线的人」——比如同一进程内多个实例同步状态。Redis 的 list / stream 能做轻量队列,但 RabbitMQ 的路由和可靠性能力它没有。 - WebSocket / SSE:服务端和浏览器之间的长连接,用于服务端主动推消息给前端。和前两个本质不同——前两个是服务和服务之间,WebSocket 是服务和浏览器之间。learnhub 的聊天室
src/modules/chat/chat.gateway.ts就是 socket.io 的WebSocketGateway,后续实时通信相关章节会展开。
判断口径:
| 需求 | 选 |
|---|---|
| 异步任务、消费者宕机不丢、要 ack/重试 | RabbitMQ(这一章) |
| 进程内广播、没订阅者就丢也无所谓 | Redis pub/sub |
| 服务端主动推浏览器的实时消息 | WebSocket / SSE |
| 超高吞吐的日志流(每秒百万级) | Kafka |
不要什么都塞 MQ——一次额外的网络跳、broker 的运维成本、消息格式约定的负担,都是代价。learnhub 选 MQ 是因为「发帖后异步推 feed / 站内信」这个需求确实需要持久化(消费者经常不在线)和解耦(生产者不该知道谁消费)。
这一章的成果
- 看懂 learnhub 的
AmqpService:@Global()模块、OnModuleInit建连、durable: true队列、persistent: true消息、try/catch 降级让 MQ 挂了不影响主流程。 - 理解 fire-and-forget 的 tradeoff:业务请求不阻塞,代价是 broker 故障窗口期会丢消息——learnhub 在「事件是副作用」场景下接受这个代价,强一致场景要上 publisher confirms 或 outbox。
- 在
PostService.create/UserService.create真实代码里看到「事务提交后发领域事件」的写法,理解事务边界画在哪。 - 自己写了一个真实的 Nest 消费者:
assertQueue对齐、prefetch控速、noAck: false+ 手动ack。 - 把旧版课程 promised 但没做的 DLX 死信队列补齐了:主队列挂
x-dead-letter-exchange、失败消息nack(false)进 DLX、DLQ 必须监控。 - 学会按需选 RabbitMQ / Redis pub/sub / WebSocket——而不是看到「传消息」就上 MQ。
常见问题
- learnhub 启动报
RabbitMQ 连接失败:正常,没起 RabbitMQ 就这样。docker compose -f docker/docker-compose.yml up -d rabbitmq起一下,或忽略——降级逻辑让主流程照样能跑。 - 生产者发了消息、消费者没收到:先看管理界面 Queues 里
learnhub.events有没有消息堆积(有堆积 = 没消费者在消费)、有没有消费者连着(Consumers 列)。如果消费者连着但消息没出,多半是 prefetch 卡住或没 ack。 - 消费者 ack 了消息还是重复消费:你看到的不是同一条消息被 ack 两次,是 broker 在你 ack 之前认为连接断了(网络抖动、消费者处理超时断开)、把消息重新投递给了另一个实例。手动 ack 之前,业务逻辑要设计成幂等的——同一条事件消费两次不该出问题(比如别用自增计数、改用 upsert)。
PRECONDITION_FAILED - inequivalent arg 'x-dead-letter-exchange':你尝试给一个已经存在、但没 DLX 属性的队列加 DLX。AMQP 不允许改属性。解决:管理界面删掉这个队列,重新声明;或者换个队列名(生产环境改名要走迁移)。- DLQ 消息一直涨怎么办:业务在持续失败。先看 DLQ 里消息的 content 和 headers(管理界面 GetMessage 能看),定位失败原因;修完后可以选择把这些消息重新 publish 回主队列(手动回放)。
- 要不要用
@nestjs/microservices的@EventPattern/@RabbitSubscribe:可以,它把这一章手写的 channel/ack 封装成了装饰器,写起来更短。代价是多一层抽象、ack 策略被框架接管,要按它的方式配置。learnhub 选了最薄的 amqplib 直连,是因为它只发不收、不需要这层抽象;如果你写复杂消费者,@nestjs/microservices值得评估。
下一章讲 Elasticsearch——全文检索引擎,解决 MySQL LIKE 搜不快、搜不准的问题。learnhub 的 SearchService 是它的真实落地,对应 PostService.create 里另一行 this.indexPost(post)。