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

NestJS中RabbitMQ多生产者与消费者配置问题咨询

NestJS 多队列RabbitMQ生产/消费配置问题

我有一个连接RabbitMQ并与Java Spring应用通信的NestJS应用,需要实现多队列的消息生产与消费,包括同一队列既生产又消费的场景。当前我的配置已支持向多队列发送消息,但不确定该实现方式是否规范;而配置多消费者时遇到问题——RmqOptions似乎仅接受单个配置对象。我的问题如下:

  1. 如何正确配置多消费者?
  2. 应如何配置多生产者?
  3. 最简洁的配置结构是什么?

现有相关文件

rabbitmq.module.ts

import { Module } from '@nestjs/common';
import { RabbitmqPublisher } from 'main/rabbit/rabbitmq.publisher';
import { ClientsModule } from '@nestjs/microservices';
import { rabbitMQPublishers } from 'main/rabbit/config/rabbitmq.publisher.config';

@Module({
  imports: [ClientsModule.register(rabbitMQPublishers)],
  providers: [RabbitmqPublisher],
  exports: [RabbitmqPublisher, ClientsModule],
})
export class RabbitMQModule {}

rabbitmq.publisher.ts

import { Inject, Injectable } from '@nestjs/common';
import { ClientProxy, RpcException } from '@nestjs/microservices';
import { RABBITMQ } from 'main/rabbit/rabbitmq.constants';

@Injectable()
export class RabbitmqPublisher {
  constructor(@Inject('CHAT') private readonly chatQueue: ClientProxy) {}

  async sendToChatQueue(data: any) {
      this.chatQueue.emit(RABBITMQ.ROUTING_KEYS.CHAT_MESSAGE, data);
  }
}

main.ts

const app = await NestFactory.create(AppModule);

app.connectMicroservice(rabbitMQConfig());
await app.startAllMicroservices();

rabbitmq.publisher.config.ts

import { ClientsModuleOptions, Transport } from '@nestjs/microservices';
import { RABBITMQ } from 'main/rabbit/rabbitmq.constants';

export const rabbitMQPublishers: ClientsModuleOptions = [
  {
    name: 'CHAT',
    transport: Transport.RMQ,
    options: {
      urls: [RABBITMQ.URL],
      queue: RABBITMQ.QUEUES.CHAT_MESSAGE_QUEUE,
      queueOptions: {
        durable: true,
        deadLetterExchange: 'dead_letter_exchange',
        maxPriority: 10,
      },
      //noAck: false,
    },
  },
];

rabbitmq.consumer.config.ts

import { RmqOptions, Transport } from '@nestjs/microservices';
import { RABBITMQ } from 'main/rabbit/rabbitmq.constants';

export const rabbitMQConfig = (): RmqOptions => ({
  transport: Transport.RMQ,
  options: {
    urls: [RABBITMQ.URL],
    queue: RABBITMQ.QUEUES.AGENT_ASSIGNMENT_QUEUE,
    queueOptions: {
      durable: true,
    },
  },
});

问题解答

1. 正确配置多消费者的方式

app.connectMicroservice()支持传入多个微服务配置,无需局限于单个RmqOptions。可以将多队列消费者配置整理成数组,循环调用connectMicroservice完成注册:

步骤1:定义多消费者配置数组

// rabbitmq.consumer.config.ts
import { RmqOptions, Transport } from '@nestjs/microservices';
import { RABBITMQ } from 'main/rabbit/rabbitmq.constants';

export const rabbitMQConsumerConfigs: RmqOptions[] = [
  {
    transport: Transport.RMQ,
    options: {
      urls: [RABBITMQ.URL],
      queue: RABBITMQ.QUEUES.AGENT_ASSIGNMENT_QUEUE,
      queueOptions: { durable: true },
    },
  },
  {
    transport: Transport.RMQ,
    options: {
      urls: [RABBITMQ.URL],
      queue: RABBITMQ.QUEUES.CHAT_MESSAGE_QUEUE, // 支持同一队列既生产又消费
      queueOptions: { durable: true },
    },
  },
];

步骤2:在main.ts中注册所有消费者

const app = await NestFactory.create(AppModule);

// 遍历配置数组,逐个连接微服务
rabbitMQConsumerConfigs.forEach(config => {
  app.connectMicroservice(config);
});

await app.startAllMicroservices();

步骤3:编写消费逻辑

在对应服务中用@EventPattern或@MessagePattern标记消费方法:

// chat.consumer.ts
@Injectable()
export class ChatConsumer {
  @EventPattern(RABBITMQ.ROUTING_KEYS.CHAT_MESSAGE)
  handleChatMessage(data: any) {
    // 处理消息逻辑
  }
}

2. 多生产者的规范配置

你当前使用ClientsModule.register()传入配置数组的方式是规范的,只需扩展配置数组并注入对应ClientProxy实例即可:

步骤1:扩展生产者配置数组

// rabbitmq.publisher.config.ts
export const rabbitMQPublishers: ClientsModuleOptions = [
  {
    name: 'CHAT',
    transport: Transport.RMQ,
    options: {
      urls: [RABBITMQ.URL],
      queue: RABBITMQ.QUEUES.CHAT_MESSAGE_QUEUE,
      queueOptions: {
        durable: true,
        deadLetterExchange: 'dead_letter_exchange',
        maxPriority: 10,
      },
    },
  },
  {
    name: 'AGENT_ASSIGNMENT',
    transport: Transport.RMQ,
    options: {
      urls: [RABBITMQ.URL],
      queue: RABBITMQ.QUEUES.AGENT_ASSIGNMENT_QUEUE,
      queueOptions: { durable: true },
    },
  },
];

步骤2:注入多生产者实例

// rabbitmq.publisher.ts
@Injectable()
export class RabbitmqPublisher {
  constructor(
    @Inject('CHAT') private readonly chatQueue: ClientProxy,
    @Inject('AGENT_ASSIGNMENT') private readonly agentQueue: ClientProxy
  ) {}

  async sendToChatQueue(data: any) {
    this.chatQueue.emit(RABBITMQ.ROUTING_KEYS.CHAT_MESSAGE, data);
  }

  async sendToAgentQueue(data: any) {
    this.agentQueue.emit(RABBITMQ.ROUTING_KEYS.AGENT_ASSIGNMENT, data);
  }
}

3. 最简洁的配置结构

将生产者、消费者配置统一管理,简化模块注册逻辑,示例如下:

统一配置文件 rabbitmq.config.ts

import { ClientsModuleOptions, RmqOptions, Transport } from '@nestjs/microservices';
import { RABBITMQ } from './rabbitmq.constants';

// 生产者配置
export const rabbitMQProducers: ClientsModuleOptions = [
  { name: 'CHAT', transport: Transport.RMQ, options: {
    urls: [RABBITMQ.URL],
    queue: RABBITMQ.QUEUES.CHAT_MESSAGE_QUEUE,
    queueOptions: { durable: true, deadLetterExchange: 'dead_letter_exchange', maxPriority: 10 }
  }},
  { name: 'AGENT_ASSIGNMENT', transport: Transport.RMQ, options: {
    urls: [RABBITMQ.URL],
    queue: RABBITMQ.QUEUES.AGENT_ASSIGNMENT_QUEUE,
    queueOptions: { durable: true }
  }}
];

// 消费者配置
export const rabbitMQConsumers: RmqOptions[] = [
  { transport: Transport.RMQ, options: {
    urls: [RABBITMQ.URL],
    queue: RABBITMQ.QUEUES.CHAT_MESSAGE_QUEUE,
    queueOptions: { durable: true }
  }},
  { transport: Transport.RMQ, options: {
    urls: [RABBITMQ.URL],
    queue: RABBITMQ.QUEUES.AGENT_ASSIGNMENT_QUEUE,
    queueOptions: { durable: true }
  }}
];

RabbitMQ模块 rabbitmq.module.ts

import { Module } from '@nestjs/common';
import { ClientsModule } from '@nestjs/microservices';
import { rabbitMQProducers } from './rabbitmq.config';
import { RabbitmqPublisher } from './rabbitmq.publisher';
import { ChatConsumer, AgentConsumer } from './consumers';

@Module({
  imports: [ClientsModule.register(rabbitMQProducers)],
  providers: [RabbitmqPublisher, ChatConsumer, AgentConsumer],
  exports: [RabbitmqPublisher],
})
export class RabbitMQModule {}

main.ts

const app = await NestFactory.create(AppModule);
import { rabbitMQConsumers } from './main/rabbit/rabbitmq.config';

rabbitMQConsumers.forEach(config => app.connectMicroservice(config));
await app.startAllMicroservices();
await app.listen(3000);

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 14:39:57