NestJS中RabbitMQ多生产者与消费者配置问题咨询
NestJS 多队列RabbitMQ生产/消费配置问题
我有一个连接RabbitMQ并与Java Spring应用通信的NestJS应用,需要实现多队列的消息生产与消费,包括同一队列既生产又消费的场景。当前我的配置已支持向多队列发送消息,但不确定该实现方式是否规范;而配置多消费者时遇到问题——RmqOptions似乎仅接受单个配置对象。我的问题如下:
- 如何正确配置多消费者?
- 应如何配置多生产者?
- 最简洁的配置结构是什么?
现有相关文件
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
相关产品推荐
相关产品推荐

