跳转到主要内容

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) 声明持久化队列、publishsendToQueue(persistent: true)
  • PostService.createUserService.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:15672guest / 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、不 awaitsendToQueue 本身是同步调用(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.createUserService.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 主仓库就这么声明过),要先在管理界面里把队列删掉再重新建。

然后第五步 catchnack(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 / 站内信」这个需求确实需要持久化(消费者经常不在线)和解耦(生产者不该知道谁消费)。

这一章的成果

  1. 看懂 learnhub 的 AmqpService@Global() 模块、OnModuleInit 建连、durable: true 队列、persistent: true 消息、try/catch 降级让 MQ 挂了不影响主流程。
  2. 理解 fire-and-forget 的 tradeoff:业务请求不阻塞,代价是 broker 故障窗口期会丢消息——learnhub 在「事件是副作用」场景下接受这个代价,强一致场景要上 publisher confirms 或 outbox。
  3. PostService.create / UserService.create 真实代码里看到「事务提交后发领域事件」的写法,理解事务边界画在哪。
  4. 自己写了一个真实的 Nest 消费者:assertQueue 对齐、prefetch 控速、noAck: false + 手动 ack
  5. 把旧版课程 promised 但没做的 DLX 死信队列补齐了:主队列挂 x-dead-letter-exchange、失败消息 nack(false) 进 DLX、DLQ 必须监控。
  6. 学会按需选 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)