Spring Boot 3.1.5中Kafka客户端无法连接远程Broker,仍访问localhost:9092
问题描述
我使用Spring Boot 3.1.5和Kafka 3.0.12搭建应用,本地Docker部署Kafka时监听器正常工作,但将Kafka迁移到另一台机器后,消费者抛出以下错误:
[org.springframework.kafka.KafkaListenerEndpointContainer#0-0-C-1] WARN o.apache.kafka.clients.NetworkClient - [Consumer clientId=consumer-testGrp-1, groupId=testGrp] Connection to node 1 (localhost/127.0.0.1:9092) could not be established. Broker may not be available.
明明在配置里指定了spring.kafka.bootstrap-servers=172.16.16.15:9092,但应用却始终尝试连接127.0.0.1,完全不使用配置的远程地址。
问题原因
这是Kafka Broker的advertised.listeners配置导致的:
客户端通过配置的bootstrap地址连接Broker后,Broker会把自己advertised.listeners配置的地址返回给客户端,后续客户端所有通信都会使用这个地址。如果Broker的advertised.listeners设置的是localhost:9092,哪怕你配置了正确的远程bootstrap地址,客户端最终还是会转向连接本地。
解决方案
1. 修改Kafka Broker配置
找到Kafka的配置文件(通常是server.properties),调整以下两项:
listeners:设置为Broker允许外部访问的监听地址,比如允许所有网卡访问:listeners=PLAINTEXT://0.0.0.0:9092advertised.listeners:设置为客户端实际能访问的Broker地址,也就是你配置在Spring Boot里的地址:advertised.listeners=PLAINTEXT://172.16.16.15:9092
如果是Docker部署的Kafka,直接在启动命令或docker-compose.yml里添加环境变量:
environment: KAFKA_LISTENERS: PLAINTEXT://0.0.0.0:9092 KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://172.16.16.15:9092
2. 重启Broker
修改配置后重启Kafka服务,确保新配置生效。
3. 验证客户端配置
确认Spring Boot应用的application.properties中spring.kafka.bootstrap-servers确实是172.16.16.15:9092,没有被环境变量、启动参数等其他配置覆盖。
用户提供的相关配置及代码
应用配置
server.port=8081 spring.kafka.bootstrap-servers=172.16.16.15:9092 spring.kafka.consumer.groupid=testGrp long.message.topic.name=longMessage greeting.topic.name=greeting filtered.topic.name=filtered partitioned.topic.name=partitioned multi.type.topic.name=multitype test.topic=testtopic1 spring.jmx.enabled=true
消费者方法
@KafkaListener(topics = "${test.topic}") public void receive(ConsumerRecord<?, ?> consumerRecord) { LOGGER.info("received payload='{}'", consumerRecord.toString()); payload = consumerRecord.toString(); latch.countDown(); }
pom.xml依赖及插件
<dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-web</artifactId> </dependency> <dependency> <groupId>org.springframework.kafka</groupId> <artifactId>spring-kafka</artifactId> </dependency> <!-- 命令行启动打包插件 --> <plugin> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-maven-plugin</artifactId> <configuration> <mainClass>com.baeldung.kafka.embedded.KafkaProducerConsumerApplication</mainClass> </configuration> <executions> <execution> <id>repackage</id> <goals> <goal>repackage</goal> </goals> </execution> </executions> </plugin>
内容的提问来源于stack exchange,提问作者Сергей Аблаев

