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

如何在单个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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 13:12:46