如何在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
相关产品推荐
相关产品推荐

