求助:如何在NestJS中正确实现多队列的RabbitMQ Fanout Exchange?
NestJS 集成 RabbitMQ Fanout Exchange 实现库存消息广播
当前实现的核心问题是未利用 Fanout Exchange 的广播特性,每个服务直接绑定独立队列,导致消息无法同时分发到邮件和日志队列。以下是修正后的完整实现方案:
1. 更新环境变量(.env)
添加 Fanout Exchange 配置:
RABBIT_MQ_URI=amqp://localhost:5672 RABBIT_MQ_STOCK_EXCHANGE=stock_update_fanout_exchange RABBIT_MQ_EMAIL_QUEUE=stock_update_email_queue RABBIT_MQ_LOG_QUEUE=stock_update_log_queue
2. 重构 RabbitMQ 配置模块
rabbitmq.module.ts
区分生产者和消费者配置,生产者绑定 Exchange,消费者绑定队列并关联到 Exchange:
import { DynamicModule, Module } from '@nestjs/common'; import { ConfigService } from '@nestjs/config'; import { ClientsModule, Transport } from '@nestjs/microservices'; import { RabbitMQService } from './rabbitmq.service'; interface RmqModuleOptions { name: string; isProducer?: boolean; } @Module({ providers: [RabbitMQService], exports: [RabbitMQService], }) export class RmqModule { static register({ name, isProducer = false }: RmqModuleOptions): DynamicModule { return { module: RmqModule, imports: [ ClientsModule.registerAsync([ { name, useFactory: (configService: ConfigService) => { const baseConfig = { transport: Transport.RMQ, options: { urls: [configService.get<string>('RABBIT_MQ_URI')], queueOptions: { durable: true }, }, }; if (isProducer) { // 生产者配置:仅指定Fanout Exchange,不绑定队列 baseConfig.options.exchange = configService.get<string>('RABBIT_MQ_STOCK_EXCHANGE'); baseConfig.options.exchangeOptions = { type: 'fanout', durable: true, }; delete baseConfig.options.queue; } else { // 消费者配置:指定队列并关联到Exchange baseConfig.options.queue = configService.get<string>(`RABBIT_MQ_${name.toUpperCase()}_QUEUE`); baseConfig.options.exchange = configService.get<string>('RABBIT_MQ_STOCK_EXCHANGE'); } return baseConfig; }, inject: [ConfigService], }, ]), ], exports: [ClientsModule], }; } }
rabbitmq.service.ts
更新配置生成逻辑,支持生产者/消费者模式切换:
import { Injectable, Logger } from '@nestjs/common'; import { ConfigService } from '@nestjs/config'; import { RmqOptions, Transport } from '@nestjs/microservices'; @Injectable() export class RabbitMQService { private readonly logger = new Logger(RabbitMQService.name); constructor(private readonly configService: ConfigService) { this.logger.log('RabbitMQService initialized'); } getOptions(queue: string, isProducer = false): RmqOptions { const options: RmqOptions = { transport: Transport.RMQ, options: { urls: [this.configService.get<string>('RABBIT_MQ_URI')], queueOptions: { durable: true }, }, }; if (isProducer) { options.options.exchange = this.configService.get<string>('RABBIT_MQ_STOCK_EXCHANGE'); options.options.exchangeOptions = { type: 'fanout', durable: true }; delete options.options.queue; } else { options.options.queue = this.configService.get<string>(`RABBIT_MQ_${queue.toUpperCase()}_QUEUE`); options.options.exchange = this.configService.get<string>('RABBIT_MQ_STOCK_EXCHANGE'); } return options; } }
3. 修改库存模块(生产者)
inventory.module.ts
注册生产者类型的 RabbitMQ 客户端:
import { Module } from '@nestjs/common'; import { InventoryService } from './inventory.service'; import { InventoryController } from './inventory.controller'; import { AccessModule } from '@app/common/access-control/access.module'; import { RedisModule } from '@app/common/redis/redis.module'; import { DatabaseModule } from 'src/database/database.module'; import { JwtService } from '@nestjs/jwt'; import { ProductModule } from 'src/product/product.module'; import { RmqModule } from '@app/common/rabbit-mq/rabbitmq.module'; @Module({ imports: [ AccessModule, RedisModule, DatabaseModule, ProductModule, RmqModule.register({ name: 'INVENTORY_PRODUCER', isProducer: true, }), ], providers: [InventoryService, JwtService], controllers: [InventoryController], }) export class InventoryModule {}
inventory.service.ts
注入生产者客户端,直接向 Fanout Exchange 发送消息:
import { Injectable, Logger, Inject } from '@nestjs/common'; import { DatabaseService } from 'src/database/database.service'; import { ProductService } from 'src/product/product.service'; import { NotFoundException, InternalServerErrorException } from '@nestjs/common'; import { ClientProxy } from '@nestjs/microservices'; @Injectable() export class InventoryService { private readonly logger = new Logger(InventoryService.name); constructor( private readonly databaseService: DatabaseService, private readonly productService: ProductService, @Inject('INVENTORY_PRODUCER') private readonly rmqProducer: ClientProxy, ) {} async updateProductStock( productId: string, quantity: number, ): Promise<any> { try { const product = await this.productService.getProductById(productId); if (!product) { throw new NotFoundException('Product not found'); } const updatedProduct = await this.databaseService.product.update({ where: { id: productId }, data: { stock: { increment: quantity } }, }); this.logger.log(`Updated product stock for productId: ${productId}, incremented by: ${quantity}`); // Fanout Exchange 忽略路由键,传入空字符串即可 this.rmqProducer.emit('', { productId, quantity }); return updatedProduct; } catch (error) { this.logger.error(`Failed to update product stock for productId: ${productId}, error: ${error.message}`); throw new InternalServerErrorException(error.message); } } }
4. 邮件/日志模块(消费者)
邮件和日志模块的代码无需大幅修改,保持原有结构即可,队列会自动绑定到 Fanout Exchange:
email.module.ts
import { Module } from '@nestjs/common'; import { RmqModule } from '@app/common/rabbit-mq/rabbitmq.module'; import { EmailService } from './email.service'; import { EmailController } from './email.controller'; @Module({ imports: [RmqModule.register({ name: 'email' })], controllers: [EmailController], providers: [EmailService], exports: [EmailService], }) export class EmailModule {}
logger.module.ts
import { Module } from '@nestjs/common'; import { RmqModule } from '@app/common/rabbit-mq/rabbitmq.module'; import { LogService } from './log.service'; import { LogController } from './log.controller'; @Module({ imports: [RmqModule.register({ name: 'logger' })], controllers: [LogController], providers: [LogService], }) export class LogModule {}
5. 更新主文件(main.ts)
分别连接生产者和消费者微服务:
import { NestFactory } from '@nestjs/core'; import { AppModule } from './app.module'; import { ValidationPipe } from '@nestjs/common'; import { Logger } from 'nestjs-pino'; import { ConfigService } from '@nestjs/config'; import * as passport from 'passport'; import * as cookieParser from 'cookie-parser'; import { RabbitMQService } from '@app/common/rabbit-mq/rabbitmq.service'; async function bootstrap() { const app = await NestFactory.create(AppModule); const rmqService = app.get<RabbitMQService>(RabbitMQService); // 连接生产者微服务 app.connectMicroservice(rmqService.getOptions('inventory', true)); // 连接邮件消费者微服务 app.connectMicroservice(rmqService.getOptions('email')); // 连接日志消费者微服务 app.connectMicroservice(rmqService.getOptions('logger')); await app.startAllMicroservices(); app.use(cookieParser()); app.use(passport.initialize()); app.useGlobalPipes( new ValidationPipe({ whitelist: true, transform: true, transformOptions: { enableImplicitConversion: true }, }), ); app.useLogger(app.get(Logger)); const configService = app.get(ConfigService); const port = configService.get('PORT'); await app.listen(port); } bootstrap();
关键说明
- Fanout Exchange 会将所有收到的消息广播到所有绑定的队列,无需路由键,生产者发送时路由键传空字符串即可。
- 消费者队列通过配置自动绑定到 Exchange,确保消息能被所有订阅服务接收。
- 生产者无需指定队列,直接将消息发送到 Exchange,由 Exchange 完成分发逻辑。
内容的提问来源于stack exchange,提问作者berk
相关产品推荐
相关产品推荐

