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

求助:如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 03:59:49