如何在单个NestJS应用中创建带不同消费者组ID的Kafka消费者?
在NestJS中创建多个不同组ID的Kafka消费者
默认通过createMicroservice启动的Kafka微服务只能配置一个全局消费者组ID。要实现多个不同组ID的消费者,你可以通过ClientsModule注册多个独立的Kafka客户端,每个客户端配置专属的消费者组ID。
步骤1:在模块中注册多个Kafka客户端
在AppModule里使用ClientsModule.register注册两个(或更多)Kafka客户端,分别指定不同的clientId和consumer.groupId:
import { Module } from '@nestjs/common'; import { ClientsModule, Transport } from '@nestjs/microservices'; import { AppService } from './app.service'; @Module({ imports: [ ClientsModule.register([ // 第一个消费者客户端,组ID为group-1 { name: 'KAFKA_GROUP_1_CLIENT', transport: Transport.KAFKA, options: { client: { clientId: 'kafka-client-1', brokers: ['localhost:9092'], }, consumer: { groupId: 'group-1', }, }, }, // 第二个消费者客户端,组ID为group-2 { name: 'KAFKA_GROUP_2_CLIENT', transport: Transport.KAFKA, options: { client: { clientId: 'kafka-client-2', brokers: ['localhost:9092'], }, consumer: { groupId: 'group-2', }, }, }, ]), ], providers: [AppService], }) export class AppModule {}
步骤2:在服务中注入并使用不同客户端
在业务服务中,通过@Inject装饰器注入不同的Kafka客户端,然后分别连接并处理消息。每个客户端会以自己的组ID消费消息:
import { Inject, Injectable, OnModuleInit } from '@nestjs/common'; import { ClientKafka, KafkaMessagePattern } from '@nestjs/microservices'; @Injectable() export class AppService implements OnModuleInit { constructor( @Inject('KAFKA_GROUP_1_CLIENT') private readonly kafkaGroup1Client: ClientKafka, @Inject('KAFKA_GROUP_2_CLIENT') private readonly kafkaGroup2Client: ClientKafka, ) {} async onModuleInit() { // 连接客户端 await this.kafkaGroup1Client.connect(); await this.kafkaGroup2Client.connect(); } // 处理group-1订阅的topic消息 @KafkaMessagePattern('target-topic') handleGroup1Message(message: any) { console.log('Group-1 收到消息:', message.value.toString()); } // 处理group-2订阅的同一topic消息(组ID不同,独立消费) @KafkaMessagePattern('target-topic') handleGroup2Message(message: any) { console.log('Group-2 收到消息:', message.value.toString()); } }
关键说明
- 每个Kafka客户端的
clientId必须唯一,避免Kafka集群识别冲突。 - 即使多个客户端订阅同一个主题,只要
groupId不同,它们会各自独立消费该主题的消息(Kafka会按组分配分区)。
内容的提问来源于stack exchange,提问作者Elgun Qasimzade
相关产品推荐
相关产品推荐

