You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Docker Compose中Kafka健康检查失败及连接超时问题求助

Kafka容器健康检查失败及服务启动问题

问题背景

我需要确保docker-compose.yml中的Kafka在我的服务启动前处于健康运行状态,该服务会向由服务端KafkaAdmin创建的Kafka Topic发送消息。

问题现象

我尝试了两种健康检查方式:获取Topic列表、查询Cluster ID。虽然日志显示Kafka已启动:

kafka-1         | [2026-05-10 11:21:31,407] INFO [KafkaRaftServer nodeId=1] Kafka Server started (kafka.server.KafkaRaftServer)

但两种健康检查均失败,报错:

dependency failed to start: container task-service-kafka-1 is unhealthy

奇怪的是,通过docker exec执行相同的健康检查命令却能正常工作:

C:\Users\nadch>docker exec -it 112c4d7c7a9362b49e3426dce609f67b1fa06233bbef65c7917753c289962924 /opt/kafka/bin/kafka-cluster.sh cluster-id --bootstrap-server localhost:9092
Cluster ID: 5L6g3nShT-eMCtK--X86sw

是否需要额外配置?根据官方文档,若未覆盖任何配置则会使用默认值,但即使粘贴文档中的环境配置也无法解决问题。

需要注意的是,若您覆盖了任何配置,则不会使用默认配置。
—— 覆盖默认Broker配置

配置文件

docker-compose.yml

version: '2.1'

services:
  task-service:
    build:
      context: .
      dockerfile: Dockerfile
    ports:
      - "8080:8080"
    depends_on:
      kafka:
        condition: service_healthy
    environment:
      KAFKA_SERVER: kafka:9092

  kafka:
    image: apache/kafka:3.9.2
    ports:
      - "9092:9092"
    # 从文档复制的配置
    environment:
      KAFKA_NODE_ID: 1
      KAFKA_PROCESS_ROLES: broker,controller
      KAFKA_LISTENERS: PLAINTEXT://localhost:9092,CONTROLLER://localhost:9093
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
      KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT
      KAFKA_CONTROLLER_QUORUM_VOTERS: 1@localhost:9093
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
      KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
      KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0
      # 改为默认值
      KAFKA_NUM_PARTITIONS: 1
    healthcheck:
#      test: ["CMD-SHELL", "kafka-topics.sh --bootstrap-server localhost:9092 --list"]
      test: ["CMD", "/kafka/bin/kafka-cluster.sh", "cluster-id", "--bootstrap-server", "localhost:9092"]
      interval: 10s
      timeout: 10s
      retries: 5
      start_period: 30s

application.yaml(仅保留相关部分)

spring:
  kafka:
    bootstrap-servers: ${KAFKA_SERVER:localhost:9092}

更新信息

即使放弃健康检查,改用service_started条件,服务启动仍失败,报错:

task-service-1  | 2026-05-10T12:01:04.299Z  INFO 1 --- [task-service] [service-admin-0] o.a.k.c.a.i.AdminMetadataManager         : [AdminClient clientId=task-service-admin-0] Rebootstrapping with Cluster(id = null, nodes = [kafka:9092 (id: -1 rack: null isFenced: false)], partitions = [], controller = null)
task-service-1  | 2026-05-10T12:01:04.366Z  INFO 1 --- [task-service] [service-admin-0] o.a.k.clients.admin.KafkaAdminClient     : [AdminClient clientId=task-service-admin-0] Forcing a hard I/O thread shutdown. Requests in progress will be aborted.
task-service-1  | 2026-05-10T12:01:04.367Z  INFO 1 --- [task-service] [service-admin-0] o.a.kafka.common.utils.AppInfoParser     : App info kafka.admin.client for task-service-admin-0 unregistered
task-service-1  | 2026-05-10T12:01:04.367Z  INFO 1 --- [task-service] [service-admin-0] o.a.k.c.a.i.AdminMetadataManager         : [AdminClient clientId=task-service-admin-0] Metadata update failed
task-service-1  |
task-service-1  | org.apache.kafka.common.errors.TimeoutException: The AdminClient thread has exited. Call: fetchMetadata

相关代码

Kafka配置类

import com.example.task_service.data.event.TaskEvent;
import com.example.task_service.data.properties.KafkaProperties;
import lombok.RequiredArgsConstructor;
import org.apache.kafka.clients.admin.AdminClientConfig;
import org.apache.kafka.clients.admin.NewTopic;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.common.serialization.UUIDSerializer;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.core.DefaultKafkaProducerFactory;
import org.springframework.kafka.core.KafkaAdmin;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.core.ProducerFactory;
import org.springframework.kafka.support.serializer.JacksonJsonSerializer;

import java.util.HashMap;
import java.util.Map;
import java.util.UUID;

@Configuration
@RequiredArgsConstructor
public class KafkaConfig {

    private final KafkaProperties properties;

    @Bean
    public KafkaAdmin kafkaAdmin() {
        Map<String, Object> configs = new HashMap<>();
        configs.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, properties.getBootstrapServers());
        return new KafkaAdmin(configs);
    }

    @Bean
    public NewTopic tasksTopic() {
        return new NewTopic("tasks", 1, (short) 1);
    }

    @Bean
    public ProducerFactory<UUID, TaskEvent> producerFactory() {
        Map<String, Object> configProps = new HashMap<>();
        configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, properties.getBootstrapServers());
        configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, UUIDSerializer.class);
        configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JacksonJsonSerializer.class);
        return new DefaultKafkaProducerFactory<>(configProps);
    }

    @Bean
    public KafkaTemplate<UUID, TaskEvent> kafkaTemplate() {
        return new KafkaTemplate<>(producerFactory());
    }
}

Kafka属性类

import lombok.Getter;
import lombok.Setter;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.stereotype.Component;

import java.util.List;

@Component
@Getter
@Setter
@ConfigurationProperties(prefix = "spring.kafka")
public class KafkaProperties {

    private List<String> bootstrapServers;
}

需要说明的是,在无Docker、本地直接运行Kafka时,应用运行完全正常。

注: 我知道K8s可能是更好的编排方式,但这只是一个Demo,所以希望尽量保持简单。

内容的提问来源于stack exchange,提问作者Sergey Zolotarev

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.01 23:37:28