Spring Kafka消费者随机断开Broker连接问题排查求助
Spring Kafka客户端重启时偶发KafkaAdmin超时与ZooKeeper连接异常问题
一、环境与配置
桌面应用每个客户端独立配置Spring Kafka生产者与消费者,核心配置如下:
Topic定义与消费者代码
@Bean public NewTopic generalTopic() { return TopicBuilder.name("topic") .partitions(5) .replicas(5) .build(); } @KafkaListener(id= "anyID", topics="topic") public void consumer(String message) { System.out.println(message); }
生产者代码
@Autowired private KafkaTemplate<String, String> kafkaTemplate; kafkaTemplate.send("topic", "message to send");
核心配置文件
spring.kafka.bootstrap-servers=localhost spring.kafka.consumer.key-deserializer=org.apache.kafka.common.serialization.StringDeserializer spring.kafka.consumer.value-deserializer=org.apache.kafka.common.serialization.StringDeserializer spring.kafka.producer.key-serializer=org.apache.kafka.common.serialization.StringSerializer spring.kafka.producer.value-serializer= org.apache.kafka.common.serialization.StringSerializer
注:@KafkaListener的id为用户名,其余配置各客户端一致。
二、问题现象
应用关闭后重启时,偶尔触发以下错误:
2022-10-21 15:13:42.130 ERROR 3524 --- [JavaFX-Launcher] o.springframework.kafka.core.KafkaAdmin : Could not configure topics org.springframework.kafka.KafkaException: Timed out waiting to get existing topics; nested exception is java.util.concurrent.TimeoutException
查看Kafka的server.log,发现关联的ZooKeeper日志:
[2022-10-21 18:11:28,519] WARN Unexpected exception (org.apache.zookeeper.server.NIOServerCnxn) EndOfStreamException: Unable to read additional data from client, it probably closed the socket: address = /127.0.0.1:34456, session = 0x100a5e36a580005 at org.apache.zookeeper.server.NIOServerCnxn.handleFailedRead(NIOServerCnxn.java:163) at org.apache.zookeeper.server.NIOServerCnxn.doIO(NIOServerCnxn.java:326) at org.apache.zookeeper.server.NIOServerCnxnFactory$IOWorkRequest.doWork(NIOServerCnxnFactory.java:522) at org.apache.zookeeper.server.WorkerService$ScheduledWorkRequest.run(WorkerService.java:154) at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) at java.lang.Thread.run(Thread.java:750) [2022-10-21 18:11:43,670] INFO Expiring session 0x100a5e36a580005, timeout of 18000ms exceeded (org.apache.zookeeper.server.ZooKeeperServer)
当前仅能通过重启Kafka服务恢复:
bin/kafka-server-start.sh -daemon config/server.properties
客户端完整错误堆栈:
org.springframework.kafka.KafkaException: Timed out waiting to get existing topics; nested exception is java.util.concurrent.TimeoutException at org.springframework.kafka.core.KafkaAdmin.lambda$checkPartitions$8(KafkaAdmin.java:388) ~[spring-kafka-2.8.8.jar:2.8.8] at java.base/java.util.HashMap.forEach(HashMap.java:1421) ~[na:na] at org.springframework.kafka.core.KafkaAdmin.checkPartitions(KafkaAdmin.java:367) ~[spring-kafka-2.8.8.jar:2.8.8] at org.springframework.kafka.core.KafkaAdmin.addOrModifyTopicsIfNeeded(KafkaAdmin.java:263) ~[spring-kafka-2.8.8.jar:2.8.8] at org.springframework.kafka.core.KafkaAdmin.initialize(KafkaAdmin.java:200) ~[spring-kafka-2.8.8.jar:2.8.8] at org.springframework.kafka.core.KafkaAdmin.afterSingletonsInstantiated(KafkaAdmin.java:167) ~[spring-kafka-2.8.8.jar:2.8.8] at org.springframework.beans.factory.support.DefaultListableBeanFactory.preInstantiateSingletons(DefaultListableBeanFactory.java:974) ~[spring-beans-5.3.22.jar:5.3.22] at org.springframework.context.support.AbstractApplicationContext.finishBeanFactoryInitialization(AbstractApplicationContext.java:918) ~[spring-context-5.3.22.jar:5.3.22] at org.springframework.context.support.AbstractApplicationContext.refresh(AbstractApplicationContext.java:583) ~[spring-context-5.3.22.jar:5.3.22] at org.springframework.boot.SpringApplication.refresh(SpringApplication.java:734) ~[spring-boot-2.7.3.jar:2.7.3] at org.springframework.boot.SpringApplication.refreshContext(SpringApplication.java:408) ~[spring-boot-2.7.3.jar:2.7.3] at org.springframework.boot.SpringApplication.run(SpringApplication.java:308) ~[spring-boot-2.7.3.jar:2.7.3] at org.springframework.boot.builder.SpringApplicationBuilder.run(SpringApplicationBuilder.java:164) ~[spring-boot-2.7.3.jar:2.7.3] at com.nume.main.MirroredMain.springBootApplicationContext(MirroredMain.java:40) ~[classes/:na] at com.nume.main.MirroredMain.init(MirroredMain.java:61) ~[classes/:na] at javafx.graphics@19/com.sun.javafx.application.LauncherImpl.launchApplication1(LauncherImpl.java:825) ~[javafx.graphics.jar:na] at javafx.graphics@19/com.sun.javafx.application.LauncherImpl.lambda$launchApplication$2(LauncherImpl.java:196) ~[javafx.graphics.jar:na] at java.base/java.lang.Thread.run(Thread.java:833) ~[na:na] Caused by: java.util.concurrent.TimeoutException: null at java.base/java.util.concurrent.CompletableFuture.timedGet(CompletableFuture.java:1960) ~[na:na] at java.base/java.util.concurrent.CompletableFuture.get(CompletableFuture.java:2095) ~[na:na] at org.apache.kafka.common.internals.KafkaFutureImpl.get(KafkaFutureImpl.java:180) ~[kafka-clients-3.1.1.jar:na] at org.springframework.kafka.core.KafkaAdmin.lambda$checkPartitions$8(KafkaAdmin.java:370) ~[spring-kafka-2.8.8.jar:2.8.8] ... 17 common frames omitted
三、排查思路
1. 核心关联分析
从日志看,客户端Socket异常关闭导致ZooKeeper会话过期,Kafka依赖ZK维护元数据,ZK会话异常会导致Kafka无法响应客户端的topic元数据查询请求,进而触发Spring KafkaAdmin的超时。
2. 调整KafkaAdmin超时配置
Spring KafkaAdmin默认超时较短,可通过配置延长超时时间,避免启动时因ZK/Kafka响应慢触发错误:
spring.kafka.admin.properties.request.timeout.ms=30000 spring.kafka.admin.properties.default.api.timeout.ms=30000
3. 优化ZooKeeper与Kafka的会话配置
- 检查ZK的
session.timeout配置,确保Kafka的zookeeper.session.timeout.ms(Kafka配置文件中)与ZK的超时匹配,避免ZK过早清理会话; - 若存在多个客户端同时连接,调整ZK的
maxClientCnxns参数,避免连接数超限导致异常关闭。
4. 确保客户端关闭时资源正常释放
桌面应用关闭时,需保证Spring容器正常销毁Kafka相关资源:
- 手动调用
KafkaListenerEndpointRegistry.stop()停止消费者容器; - 调用
kafkaTemplate.destroy()释放生产者资源,确保ZK会话正常注销,避免ZK残留无效会话。
5. 修正Topic副本数配置
当前Topic设置replicas=5,若为单Broker Kafka集群,该配置不合理(副本数不能大于Broker节点数),会导致Kafka在检查topic时持续等待,最终触发超时。单Broker环境需修改为:
@Bean public NewTopic generalTopic() { return TopicBuilder.name("topic") .partitions(5) .replicas(1) .build(); }
内容的提问来源于stack exchange,提问作者music bot
相关产品推荐
相关产品推荐

