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

多实例REST API基于Kafka实现生产者/消费者同步H2数据库问题排查

问题描述

我有一个集成H2内存数据库的简易REST API,计划启动该应用的多个实例,每个实例拥有独立的内存数据库,需要同步这些数据库。选择Kafka作为解决方案:当8080端口的实例收到POST请求时,其他所有实例也应同步执行该操作。目前每个应用实例同时作为生产者和消费者,但发送消息后只有一个实例能接收到消息。相关代码如下:

生产者配置类

@EnableKafka
@Configuration
public class KafkaProducerConfigForDepartment {

    @Value(value = "${kafka.bootstrapAddress}")
    private String bootstrapAddress;

    @Bean
    public ProducerFactory<String, MessageEventForDepartment> producerFactoryForDepartment() {
        Map<String, Object> configProps = new HashMap<>();
        configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapAddress);
        configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class);
        return new DefaultKafkaProducerFactory<>(configProps);
    }

    @Bean
    public KafkaTemplate<String, MessageEventForDepartment> kafkaTemplate() {
        return new KafkaTemplate<>(producerFactoryForDepartment());
    }

}

消费者与Topic配置类

@Configuration
public class KafkaTopicConfig {

    @Value(value = "${kafka.bootstrapAddress}")
    private String bootstrapAddress;

    @Bean
    public ConsumerFactory<String, MessageEventForDepartment> consumerFactoryForDepartments() {
        Map<String, Object> props = new HashMap<>();
        props.put(JsonDeserializer.TRUSTED_PACKAGES, "*");
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapAddress);
        props.put(ConsumerConfig.GROUP_ID_CONFIG, "groupId");
        return new DefaultKafkaConsumerFactory<>(props, new StringDeserializer(), new JsonDeserializer<>(MessageEventForDepartment.class));
    }

    @Bean
    public NewTopic topic1() {
        return TopicBuilder.name("topic12")
                .partitions(10)
                .replicas(10)
                .build();
    }

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, MessageEventForDepartment>
    kafkaListenerContainerFactoryForDepartments() {

        ConcurrentKafkaListenerContainerFactory<String, MessageEventForDepartment> factory =
                new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactoryForDepartments());
        return factory;
    }

}

消费者监听类

@Component
@Slf4j
public class DepartmentKafkaService {

    @Autowired
    private DepartmentService departmentService;

    @KafkaListener(topics = "topic12" , groupId = "groupId",containerFactory = "kafkaListenerContainerFactoryForDepartments")
    public void listenGroupFoo(MessageEventForDepartment message) {
         log.info(message.toString());
    }
}

请问这是什么原因导致的?或者我的方案是否存在问题,有什么建议?


问题原因与解决方案

核心原因

所有应用实例都使用了相同的消费者组ID(groupId)。Kafka的消费者组机制是:一个Topic的每个分区,只能被同一个消费者组内的一个消费者消费。当多个实例同属一个组时,消息会被分摊到组内的不同实例,而非广播给所有实例,这就导致你看到的“只有一个实例能收到消息”的现象。

解决方案

要实现所有实例都同步接收消息,需要让每个实例成为独立的消费者组,具体可以这样做:

  1. 为每个实例配置唯一的消费者组ID

    • 避免硬编码groupId,改为通过配置文件动态生成,比如利用实例的端口号、随机字符串或者容器ID作为组ID的一部分,确保每个实例的groupId唯一。
    • 修改配置类,从配置文件读取groupId:
      @Configuration
      public class KafkaTopicConfig {
      
          @Value(value = "${kafka.bootstrapAddress}")
          private String bootstrapAddress;
          @Value("${kafka.consumer.group-id}")
          private String groupId;
      
          @Bean
          public ConsumerFactory<String, MessageEventForDepartment> consumerFactoryForDepartments() {
              Map<String, Object> props = new HashMap<>();
              props.put(JsonDeserializer.TRUSTED_PACKAGES, "*");
              props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapAddress);
              props.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
              return new DefaultKafkaConsumerFactory<>(props, new StringDeserializer(), new JsonDeserializer<>(MessageEventForDepartment.class));
          }
      
          // 其他Bean保持不变
      }
      
    • 在application.properties或application.yml中配置:
      kafka.consumer.group-id=department-group-${server.port}
      
    • 同时修改@KafkaListener注解中的groupId,同样使用配置值:
      @KafkaListener(topics = "topic12" , groupId = "${kafka.consumer.group-id}",containerFactory = "kafkaListenerContainerFactoryForDepartments")
      public void listenGroupFoo(MessageEventForDepartment message) {
           log.info(message.toString());
      }
      
  2. 额外注意事项

    • 如果需要避免实例消费自己发送的消息,可以在发送的MessageEventForDepartment中添加实例标识(比如端口号),消费时判断消息是否由当前实例发送,若是则跳过处理。
    • 你的Topic配置了10个分区和10个副本,这个配置没问题,但要确保Kafka集群的节点数足够支撑10个副本(副本数不能超过集群节点数),否则会导致Topic创建失败。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 14:25:23