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

响应式Kafka项目配置多套消费者仅单个启动问题咨询

解决方案

同一Spring项目中完全可以同时运行多个Kafka消费者,你当前只有单个消费者启动是配置和代码中的多个错误导致的,按以下步骤调整即可实现双消费者正常运行:

核心问题排查&修正

1. 配置类未加载

确认两个Kafka配置类都添加了@Configuration注解,且所在包路径在Spring组件扫描范围内,否则Service2的相关Bean不会被初始化,自然无法启动对应的消费者。

2. 配置键拼写错误

你的application.yml中定义的源主题配置是spring.kafka.sourcetopic2,但Service2KafkaConfig中注入的是${spring.kafka.sourcetopic3},配置不存在会直接导致Bean初始化失败,需要修正为:

@Bean("service2KafkaReceiverOptions")
public ReceiverOptions<String, String> service2KafkaReceiverOptions(@Value("${spring.kafka.sourcetopic2}") String topic, KafkaProperties kafkaProperties) {
    // 原有逻辑不变
}

3. Bean名称不匹配

Spring默认以方法名作为Bean的注册名称:

  • Service1中注入的Bean名称是service1KafkaConsumerTemplate,但KafkaConfig中定义的Bean方法名为kafkaConsumerTemplate,如果没有加@Primary或者@Qualifier匹配,要么注入失败,要么只会优先加载第一个匹配的Bean。
    修正方案:给所有Bean指定明确名称,注入时通过@Qualifier匹配:

Service1配置调整

@Configuration
public class KafkaConfig {
     @Bean("service1KafkaReceiverOptions")
     public ReceiverOptions<String, String> kafkaReceiverOptions(@Value("${spring.kafka.sourcetopic1}") String topic, 
                                                                 @Value("${spring.kafka.consumer.service1-group-id}") String groupId,
                                                                 @Value("${spring.kafka.consumer.service1-client-id}") String clientId,
                                                                 KafkaProperties kafkaProperties) {
         Map<String, Object> consumerProps = kafkaProperties.buildConsumerProperties();
         // 单独设置Service1的消费者分组、客户端ID,避免和Service2冲突
         consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
         consumerProps.put(ConsumerConfig.CLIENT_ID_CONFIG, clientId);
         ReceiverOptions<String, String> basicReceiverOptions = ReceiverOptions.create(consumerProps);
         return basicReceiverOptions.subscription(Collections.singletonList(topic));
     }  

     @Bean("service1KafkaConsumerTemplate")
     public ReactiveKafkaConsumerTemplate<String, String> kafkaConsumerTemplate(@Qualifier("service1KafkaReceiverOptions") ReceiverOptions<String, String> kafkaReceiverOptions) {
         return new ReactiveKafkaConsumerTemplate<String, String>(kafkaReceiverOptions);
     }  
}

Service2配置调整

@Configuration
public class Service2KafkaConfig {
     @Bean("service2KafkaReceiverOptions")
     public ReceiverOptions<String, String> service2KafkaReceiverOptions(@Value("${spring.kafka.sourcetopic2}") String topic, 
                                                                         @Value("${spring.kafka.consumer.service2-group-id}") String groupId,
                                                                         @Value("${spring.kafka.consumer.service2-client-id}") String clientId,
                                                                         KafkaProperties kafkaProperties) {
         Map<String, Object> consumerProps = kafkaProperties.buildConsumerProperties();
         consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
         consumerProps.put(ConsumerConfig.CLIENT_ID_CONFIG, clientId);
         ReceiverOptions<String, String> basicReceiverOptions = ReceiverOptions.create(consumerProps);
         return basicReceiverOptions.subscription(Collections.singletonList(topic));
     }

     @Bean("service2KafkaConsumerTemplate")
     public ReactiveKafkaConsumerTemplate<String, String> service2KafkaConsumerTemplate(@Qualifier("service2KafkaReceiverOptions") ReceiverOptions<String, String> kafkaReceiverOptions) {
         return new ReactiveKafkaConsumerTemplate<String, String>(kafkaReceiverOptions);
     }
}

4. YML配置补充拆分消费者参数

在原有yml基础上新增两个消费者的独立分组、客户端ID配置,避免消费重平衡异常:

spring:
  kafka:
    bootstrap-servers: localhost:9092
    sourcetopic1: topic1
    destinationtopic1 : topic2
    sourcetopic2: topic3
    destinationtopic2 : topic4
    producer:
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
    consumer:
      group-id: TestCollector
      client-id:TestCollector01
      # 新增独立消费者配置
      service1-group-id: TestCollector-service1
      service1-client-id: TestCollector01-service1
      service2-group-id: TestCollector-service2
      service2-client-id: TestCollector01-service2
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.apache.kafka.common.serialization.StringDeserializer

5. 代码笔误修正

Service1代码中的.donOnNext(s->ProcessRecord(s))拼写错误,需要改为.doOnNext(s->ProcessRecord(s)),否则会导致运行时异常。

调整完成后重启项目,两个消费者会分别订阅topic1和topic3,独立运行互不影响。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.05 03:36:01