Docker部署Kafka与Spring Boot应用无法连接Broker问题求助
问题描述
尝试用Docker Compose运行一个与Kafka通信的Spring Boot应用,启动后无法连接到Broker,循环出现以下日志:
INFO 1 org.apache.kafka.clients.NetworkClient : [AdminClient clientId=adminclient-1] Node -1 disconnected. WARN 1 org.apache.kafka.clients.NetworkClient : [AdminClient clientId=adminclient-1] Connection to node -1 (localhost/127.0.0.1:9092) could not be established. Broker may not be available.
相关配置
docker-compose.yml
version: '3' networks: kafka-net: name: kafka-net driver: bridge services: zookeeper: image: confluentinc/cp-zookeeper:latest environment: ZOOKEEPER_CLIENT_PORT: 2181 networks: - kafka-net ports: - 22181:2181 kafka: image: confluentinc/cp-kafka:latest environment: KAFKA_ADVERTISED_HOST_NAME: 127.0.0.1 KAFKA_LISTENERS: DOCKER_INTERNAL://:29092,DOCKER_EXTERNAL://:9092 KAFKA_ADVERTISED_LISTENERS: DOCKER_INTERNAL://kafka:29092,DOCKER_EXTERNAL://${DOCKER_HOST_IP:-127.0.0.1}:9092 KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: DOCKER_INTERNAL:PLAINTEXT,DOCKER_EXTERNAL:PLAINTEXT KAFKA_INTER_BROKER_LISTENER_NAME: DOCKER_INTERNAL KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 KAFKA_BROKER_ID: 1 KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 KAFKA_AUTO_CREATE_TOPICS_ENABLE: true networks: - kafka-net depends_on: - zookeeper ports: - 9092:9092 common-kafka: build: context: ./common-kafka environment: KAFKA_BROKERCONNECT: kafka:29092 networks: - kafka-net depends_on: - kafka ports: - 8081:8081
application.properties
server.port=8081 kafka.bootstrap-servers=127.0.0.1:9092 kafka.producer.retries=5 kafka.transfer.topic=TRANSFER kafka.transfer.topic.dlt=TRANSFER-dlt kafka.transfer.log.group=log_transfer_group kafka.transfer.email.group=send_email_transfer_group
Kafka配置类
@Configuration public class ConfigKafkaConsumer { @Value(value = "${kafka.bootstrap-servers}") private String BOOTSTRAP_ADDRESS; @Bean public ConsumerFactory<String, TransferNotification> consumerFactory() { return new DefaultKafkaConsumerFactory<>(properties()); } @Bean public ConcurrentKafkaListenerContainerFactory<String, TransferNotification> kafkaListenerContainerFactory() { ConcurrentKafkaListenerContainerFactory<String, TransferNotification> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); return factory; } private Map<String,Object> properties() { Map<String, Object> configProps = new HashMap<>(); configProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, BOOTSTRAP_ADDRESS); configProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); configProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class); configProps.put(JsonDeserializer.TRUSTED_PACKAGES, "br.com.bank.*"); configProps.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 10); configProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); return configProps; } } @Configuration public class ConfigKafkaProducer { @Value(value = "${kafka.bootstrap-servers}") private String BOOTSTRAP_ADDRESS; @Bean public ProducerFactory<String, TransferNotification> producerFactory() { return new DefaultKafkaProducerFactory<>(properties()); } @Bean public KafkaTemplate<String, TransferNotification> kafkaTemplate() { return new KafkaTemplate<>(producerFactory()); } private Map<String, Object> properties() { Map<String, Object> configProps = new HashMap<>(); configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, BOOTSTRAP_ADDRESS); configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class); configProps.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, "true"); return configProps; } }
Kafka监听器代码
@KafkaListener(topics = "${kafka.transfer.topic}", groupId = "${kafka.transfer.email.group}") @RetryableTopic(attempts = "${kafka.producer.retries}", backoff = @Backoff(30000L)) @KafkaListener(topics = "${kafka.transfer.topic.dlt}", groupId = "${kafka.transfer.email.group}") @KafkaListener(topics = "${kafka.transfer.topic}", groupId = "${kafka.transfer.log.group}")
请问有什么解决该问题的思路?
解决思路
1. 修正Spring Boot的Kafka地址配置
Spring Boot容器和Kafka容器处于同一个Docker网络kafka-net中,容器内的localhost指向自身而非Kafka容器。你已经在docker-compose.yml中给Spring Boot服务配置了正确的环境变量KAFKA_BROKERCONNECT: kafka:29092,但application.properties里仍用了错误的127.0.0.1:9092,这是核心问题。
修改application.properties:
kafka.bootstrap-servers=kafka:29092
或者配置成兼容本地开发和Docker部署的形式:
kafka.bootstrap-servers=${KAFKA_BROKERCONNECT:127.0.0.1:9092}
2. 清理Kafka冗余配置(推荐)
KAFKA_ADVERTISED_HOST_NAME参数在新版Confluent Kafka中已被弃用,且你已经配置了KAFKA_ADVERTISED_LISTENERS,保留该参数可能导致配置冲突,建议从Kafka服务的环境变量中移除它。
3. 确保Kafka就绪后再启动应用
depends_on仅保证容器启动顺序,不保证Kafka服务完全就绪。可以通过以下方式避免应用提前启动:
- 在Spring Boot中添加Kafka健康检查逻辑,等待Broker可用后再初始化消费者/生产者
- 在Docker Compose中使用
wait-for-it脚本,让Spring Boot容器等待Kafka的29092端口就绪后再启动,示例:
common-kafka: build: context: ./common-kafka environment: KAFKA_BROKERCONNECT: kafka:29092 networks: - kafka-net depends_on: - kafka ports: - 8081:8081 command: ["./wait-for-it.sh", "kafka:29092", "--", "java", "-jar", "你的应用jar包名称.jar"]
需要将wait-for-it.sh放入应用镜像中,或使用内置该脚本的基础镜像。
4. 排查网络连通性
进入Spring Boot容器内,执行以下命令测试与Kafka容器的连通性,排除网络问题:
# 进入容器 docker exec -it <你的common-kafka容器ID> /bin/bash # 测试端口连通 telnet kafka 29092
内容的提问来源于stack exchange,提问作者Victor Soares
相关产品推荐
相关产品推荐

