多实例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的每个分区,只能被同一个消费者组内的一个消费者消费。当多个实例同属一个组时,消息会被分摊到组内的不同实例,而非广播给所有实例,这就导致你看到的“只有一个实例能收到消息”的现象。
解决方案
要实现所有实例都同步接收消息,需要让每个实例成为独立的消费者组,具体可以这样做:
为每个实例配置唯一的消费者组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()); }
额外注意事项
- 如果需要避免实例消费自己发送的消息,可以在发送的
MessageEventForDepartment中添加实例标识(比如端口号),消费时判断消息是否由当前实例发送,若是则跳过处理。 - 你的Topic配置了10个分区和10个副本,这个配置没问题,但要确保Kafka集群的节点数足够支撑10个副本(副本数不能超过集群节点数),否则会导致Topic创建失败。
- 如果需要避免实例消费自己发送的消息,可以在发送的
内容的提问来源于stack exchange,提问作者Renzo Redoran
相关产品推荐
相关产品推荐

