如何在同一Nest.js项目中用RabbitMQ作消息Broker、Kafka作事件总线?
在NestJS中同时使用RabbitMQ和Kafka的解决方案
NestJS 并不限制你只能使用一种消息传输协议,完全可以同时集成RabbitMQ和Kafka来满足你的架构需求,以下是几种可行方案:
方案1:启动多个独立的微服务实例
分别创建并启动两个微服务实例,一个绑定RabbitMQ负责网关与微服务间的通信,另一个绑定Kafka处理事件发布订阅:
import { NestFactory } from '@nestjs/core'; import { Transport, MicroserviceOptions } from '@nestjs/microservices'; import { AppModule } from './app.module'; async function bootstrap() { // 启动RabbitMQ微服务,处理网关与其他微服务的通信 const rmqApp = await NestFactory.createMicroservice<MicroserviceOptions>(AppModule, { transport: Transport.RMQ, options: { urls: ['amqp://localhost:5672'], queue: 'api_gateway_queue', queueOptions: { durable: false, }, }, }); // 启动Kafka微服务,处理UserCreated这类事件的发布与订阅 const kafkaApp = await NestFactory.createMicroservice<MicroserviceOptions>(AppModule, { transport: Transport.KAFKA, options: { client: { brokers: ['localhost:9092'], }, consumer: { groupId: 'event-handler-consumer-group', }, }, }); // 同时启动两个微服务 await Promise.all([rmqApp.listen(), kafkaApp.listen()]); } bootstrap();
方案2:在主应用中连接多个微服务传输器
如果网关本身需要暴露HTTP接口,可以在主HTTP应用中同时连接RabbitMQ和Kafka两种微服务传输器:
import { NestFactory } from '@nestjs/core'; import { Transport, MicroserviceOptions } from '@nestjs/microservices'; import { AppModule } from './app.module'; async function bootstrap() { // 创建主HTTP应用 const app = await NestFactory.create(AppModule); // 连接RabbitMQ微服务 const rmqMicroservice = app.connectMicroservice<MicroserviceOptions>({ transport: Transport.RMQ, options: { urls: ['amqp://localhost:5672'], queue: 'api_gateway_queue', queueOptions: { durable: false }, }, }); // 连接Kafka微服务 const kafkaMicroservice = app.connectMicroservice<MicroserviceOptions>({ transport: Transport.KAFKA, options: { client: { brokers: ['localhost:9092'] }, consumer: { groupId: 'gateway-kafka-consumer' }, }, }); // 启动所有微服务和主HTTP应用 await Promise.all([rmqMicroservice.listen(), kafkaMicroservice.listen()]); await app.listen(3000); // 网关的HTTP端口 } bootstrap();
方案3:通过客户端注册使用另一种消息系统
如果只需要一种作为微服务主传输器,另一种仅用于事件发布/订阅,可以在模块中通过ClientsModule注册对应的客户端:
1. 在模块中注册Kafka客户端
import { Module } from '@nestjs/common'; import { ClientsModule, Transport } from '@nestjs/microservices'; import { AppController } from './app.controller'; import { AppService } from './app.service'; @Module({ imports: [ ClientsModule.register([ { name: 'KAFKA_EVENT_CLIENT', transport: Transport.KAFKA, options: { client: { brokers: ['localhost:9092'] }, consumer: { groupId: 'user-event-consumer' }, }, }, ]), ], controllers: [AppController], providers: [AppService], }) export class AppModule {}
2. 在服务中注入并使用Kafka客户端
import { Injectable, Inject } from '@nestjs/common'; import { ClientKafka } from '@nestjs/microservices'; @Injectable() export class AppService { constructor( @Inject('KAFKA_EVENT_CLIENT') private readonly kafkaClient: ClientKafka ) {} // 发布UserCreated事件 async publishUserCreatedEvent(userData: any) { await this.kafkaClient.emit('UserCreated', userData).toPromise(); } }
这种方式下,微服务既能通过RabbitMQ处理网关通信,又能通过Kafka客户端完成事件的发布与订阅。
内容的提问来源于stack exchange,提问作者ögeday öztoprak
相关产品推荐
相关产品推荐

