动态注册/注销Kafka主题引发的偏移量同步异常
Spring Kafka动态灾备切换后MM2偏移量同步异常排查与解决
问题背景
基于动态Kafka主题监听实现消费者灾备机制,切换回主集群后,MM2无法将主集群的偏移量同步到备集群的镜像主题。控制台消费者/原生Java客户端无此问题,疑似Spring Kafka动态注销主题后存在资源泄漏。
环境版本:
- Kafka: kafka_2.13-3.1.0
- Spring Kafka: 2.9.6
灾备流程:
- 正常消费集群A的
t2主题 - 灾备触发,切换消费集群B的
A.t2镜像主题并提交偏移量 - 集群A恢复,通过MM2双向同步将B的偏移量同步到A
- 注销集群B的
A.t2主题监听,切回集群A消费t2 - 后续消费集群A
t2时,偏移量无法同步到集群B的A.t2
核心原因分析
Spring Kafka的ConcurrentMessageListenerContainer在动态注销主题时,若仅停止容器未完全销毁,会残留底层KafkaConsumer实例或相关资源:
- 残留的集群B消费者仍持有消费者组G1的连接,干扰MM2的偏移量同步逻辑
- 未释放的消费者实例会继续向集群B发送心跳或提交偏移量,覆盖MM2同步的偏移量
解决方案
1. 动态注销时完全销毁容器资源
管理监听容器时,需调用stop()后执行destroy(),确保彻底释放底层资源:
// 示例:获取目标容器后执行销毁操作 ConcurrentMessageListenerContainer<String, String> container = ...; container.stop(); container.destroy(); // 关键步骤:销毁容器,释放KafkaConsumer实例
2. 为集群A/B配置独立的消费者工厂
避免单工厂切换集群时的配置残留,为两个集群分别创建独立的消费者工厂和容器工厂:
// 集群A消费者工厂 @Bean(name = "consumerFactoryA") public ConsumerFactory<String, String> consumerFactoryA() { Map<String, Object> configProps = new HashMap<>(); configProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9093,localhost:9094"); configProps.put(ConsumerConfig.GROUP_ID_CONFIG, "G1"); configProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); configProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); configProps.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); configProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest"); return new DefaultKafkaConsumerFactory<>(configProps); } // 集群B消费者工厂 @Bean(name = "consumerFactoryB") public ConsumerFactory<String, String> consumerFactoryB() { Map<String, Object> configProps = new HashMap<>(); configProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9095,localhost:9096"); configProps.put(ConsumerConfig.GROUP_ID_CONFIG, "G1"); configProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); configProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); configProps.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); configProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest"); return new DefaultKafkaConsumerFactory<>(configProps); } // 集群A容器工厂 @Bean(name = "kafkaListenerContainerFactoryA") public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactoryA() { ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactoryA()); factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE); factory.setConcurrency(1); return factory; } // 集群B容器工厂 @Bean(name = "kafkaListenerContainerFactoryB") public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactoryB() { ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactoryB()); factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE); factory.setConcurrency(1); return factory; }
3. 优化MM2配置,排除无效消费者干扰
调整MM2配置,增强组同步的准确性:
# 严格过滤无效消费者组 groups.blacklist=console-consumer-.*, connect-.*,__.*, G1-temp-.* # 缩短组元数据刷新和偏移量同步间隔 refresh.groups.interval.seconds=30 sync.group.offsets.interval.seconds=10
4. 切换集群后手动触发MM2偏移量同步
切回集群A后,手动执行MM2偏移量同步命令,确保初始同步完成:
# 同步集群A的G1组偏移量到集群B的A.t2主题 kafka-mirror-maker.sh --consumer.config config/A-consumer.properties --producer.config config/B-producer.properties --whitelist "t2" --sync-group-offsets G1
验证步骤
- 切换回集群A前,检查JVM线程,确认无残留的KafkaConsumer线程
- 消费集群A
t2并提交偏移量后,查看集群B__consumer_offsets中G1组对应A.t2的偏移量是否更新 - 查看MM2日志,确认偏移量同步任务无报错
内容的提问来源于stack exchange,提问作者Tech_Dummy
相关产品推荐
相关产品推荐

