You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在NestJS中实现RabbitMQ Fanout交换机?及包使用故障排查

排查与解决步骤

1. 生产者发送消息改为异步非阻塞

golevelup的RabbitMQ生产者如果用同步方式发送消息,会直接阻塞Node.js事件循环,导致HTTP请求无法处理。必须采用异步发送,且不要在HTTP请求链路中同步等待消息确认。

示例代码:

// InventoryService 中的库存更新逻辑
import { RabbitMQService } from '@golevelup/nestjs-rabbitmq';

@Injectable()
export class InventoryService {
  constructor(private readonly rabbitmqService: RabbitMQService) {}

  async updateStock(productId: string, quantity: number) {
    // 先完成数据库库存更新
    const updatedStock = await this.stockRepository.update(productId, { quantity });
    
    // 异步发送消息,不阻塞HTTP响应流程
    this.rabbitmqService.publish(
      'stock_updates_exchange',
      '', // Fanout交换机无需路由键
      { productId, quantity, timestamp: new Date().toISOString() },
      { persistent: true } // 可选:确保消息持久化
    ).catch(err => {
      // 单独处理发送失败,不影响HTTP响应
      console.error('库存更新消息发送失败:', err);
    });

    // 直接返回响应,无需等待消息发送完成
    return updatedStock;
  }
}

2. 消费者逻辑必须异步,避免阻塞事件循环

EmailService和LogService的消费逻辑如果包含同步IO操作(比如同步发邮件、同步写文件),会死死卡住事件循环。所有IO操作必须改成异步方式,同时做好消息确认/拒绝处理。

示例消费者代码:

// EmailService 消息消费逻辑
import { RabbitSubscribe } from '@golevelup/nestjs-rabbitmq';

@Injectable()
export class EmailService {
  constructor(private readonly mailerService: MailerService) {}

  @RabbitSubscribe({
    exchange: 'stock_updates_exchange',
    routingKey: '',
    queue: 'email_queue',
    queueOptions: { durable: true },
  })
  async handleStockUpdate(message: { productId: string, quantity: number }) {
    try {
      // 异步发送邮件
      await this.mailerService.sendMail({
        to: 'admin@example.com',
        subject: `商品${message.productId}库存更新`,
        text: `最新库存:${message.quantity}`,
      });
      // 手动确认消息消费完成
      (message as any).ack();
    } catch (err) {
      console.error('库存更新邮件发送失败:', err);
      // 消费失败时拒绝消息,避免重复消费(根据需求调整是否重新入队)
      (message as any).nack(false, false);
    }
  }
}

// LogService 消费逻辑同理
@Injectable()
export class LogService {
  constructor(private readonly logRepository: LogRepository) {}

  @RabbitSubscribe({
    exchange: 'stock_updates_exchange',
    routingKey: '',
    queue: 'log_queue',
    queueOptions: { durable: true },
  })
  async handleStockUpdate(message: { productId: string, quantity: number, timestamp: string }) {
    try {
      // 异步写入日志到数据库
      await this.logRepository.create({
        event: 'STOCK_UPDATE',
        payload: message,
        timestamp: new Date(message.timestamp),
      });
      (message as any).ack();
    } catch (err) {
      console.error('库存更新日志写入失败:', err);
      (message as any).nack(false, false);
    }
  }
}

3. 检查RabbitMQ模块配置,确保连接非阻塞

在AppModule中配置RabbitMQ时,要避免等待连接建立完成再启动应用,同时限制消费者预取消息数量,防止一次性接收过多消息导致阻塞。

示例配置:

import { RabbitMQModule } from '@golevelup/nestjs-rabbitmq';

@Module({
  imports: [
    RabbitMQModule.forRootAsync(RabbitMQModule, {
      useFactory: () => ({
        exchanges: [
          {
            name: 'stock_updates_exchange',
            type: 'fanout',
            durable: true, // 持久化交换机
          },
        ],
        uri: 'amqp://localhost:5672',
        connectionInitOptions: {
          wait: false, // 不等待连接建立完成再启动应用
          reject: false, // 连接失败时不终止应用,自动重连
        },
        defaultConsumerOptions: {
          prefetchCount: 5, // 限制消费者同时处理的消息数量
        },
      }),
    }),
    // 其他业务模块
  ],
})
export class AppModule {}

4. 排查事件循环阻塞点

如果以上配置都没问题,用Node.js自带工具排查阻塞:

  • 启动应用时添加参数:node --trace-event-categories v8,node.async_hooks dist/main.js,查看事件循环阻塞的具体位置
  • 确保所有数据库操作、外部API调用都是异步的,禁止使用同步方法(比如fs.readFileSync)

内容的提问来源于stack exchange,提问作者berk

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.22 00:22:16