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

Docker环境下Service A连接Kafka报NoBrokersAvailable错误求助

Docker中Service A连接Kafka报错NoBrokersAvailable

我在Docker中部署了多个服务,包括作为消费者的Service A和Kafka服务,相关配置及代码如下:

docker-compose.yaml配置

version: '3.8'

services:
  service-a:
    container_name: service-a
    build:
      context: .
      dockerfile: Dockerfile
    ports:
      - "8881:8881"
    environment:
      - ENVIRONMENT=development
    depends_on:
      - kafka
      - chroma

  kafka:
    container_name: my-kafka
    image: bitnami/kafka:latest
    ports:
      - "9094:9094"
    environment:
      - KAFKA_ENABLE_KRAFT=yes
      - KAFKA_CFG_BROKER_ID=1
      - KAFKA_CFG_NODE_ID=1
      - KAFKA_CFG_PROCESS_ROLES=broker,controller
      - KAFKA_CFG_CONTROLLER_LISTENER_NAMES=CONTROLLER
      - KAFKA_CFG_LISTENERS=PLAINTEXT://:9092,CONTROLLER://:9093,EXTERNAL://:9094
      - KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT,EXTERNAL:PLAINTEXT
      - KAFKA_CFG_ADVERTISED_LISTENERS=PLAINTEXT://kafka:9092,EXTERNAL://localhost:9094
      - KAFKA_CFG_CONTROLLER_QUORUM_VOTERS=1@:9093
      - ALLOW_PLAINTEXT_LISTENER=yes    

  chroma:
    container_name: learnsuite-chroma
    image: chromadb/chroma:latest  
    ports:
      - "8882:8000"

Service A的Kafka配置

bootstrap_servers: kafka:9092

Service A消费者创建逻辑

def get_consumer_for_topic(topic: str):
    group_id = settings.kafka.group
    bootstrap = settings.kafka.bootstrap_servers
    bootstrap_servers = bootstrap.split(',') if bootstrap else []
    admin_client = KafkaAdminClient(bootstrap_servers=bootstrap_servers)

    try:
        topics = admin_client.list_topics()
        if topic not in topics:
            # Create the topic if it doesn't exist
            new_topic = NewTopic(name=topic, num_partitions=1, replication_factor=1)
            admin_client.create_topics([new_topic])
            logger.info(f"Topic '{topic}' created.")
        else:
            logger.info(f"Topic '{topic}' already exists.")
    except Exception as e:
        logger.error(f"Error ensuring topic exists: {e}")
    finally:
        admin_client.close() 
    
    return KafkaConsumer(
        topic,
        bootstrap_servers=bootstrap_servers,
        group_id = group_id,
        auto_offset_reset='earliest'
    )

执行docker-compose up后,出现以下错误:

kafka.errors.NoBrokersAvailable: NoBrokersAvailable


解决方案

问题根源是Service A启动速度过快,在Kafka完全初始化完成前就执行了主题创建操作。添加重试等待机制即可解决,修改后的消费者创建逻辑如下:

import time
from kafka.errors import NoBrokersAvailable

def get_consumer_for_topic(topic: str):
    group_id = settings.kafka.group
    bootstrap = settings.kafka.bootstrap_servers
    bootstrap_servers = bootstrap.split(',') if bootstrap else []
    max_retries = 5
    retry_delay = 3  # 每次重试间隔3秒
    admin_client = None

    for attempt in range(max_retries):
        try:
            admin_client = KafkaAdminClient(bootstrap_servers=bootstrap_servers)
            topics = admin_client.list_topics()
            
            if topic not in topics:
                new_topic = NewTopic(name=topic, num_partitions=1, replication_factor=1)
                admin_client.create_topics([new_topic])
                logger.info(f"Topic '{topic}' created.")
            else:
                logger.info(f"Topic '{topic}' already exists.")
            break  # 成功则跳出循环
        except NoBrokersAvailable:
            logger.warning(f"Attempt {attempt+1}/{max_retries}: Kafka broker not available, retrying in {retry_delay}s...")
            time.sleep(retry_delay)
        except Exception as e:
            logger.error(f"Error ensuring topic exists: {e}")
            raise
        finally:
            if admin_client:
                admin_client.close()
    else:
        # 所有重试都失败
        raise RuntimeError(f"Failed to connect to Kafka after {max_retries} attempts")
    
    return KafkaConsumer(
        topic,
        bootstrap_servers=bootstrap_servers,
        group_id = group_id,
        auto_offset_reset='earliest'
    )

修改后,Service A会在Kafka未就绪时自动重试连接,直到Kafka完全启动或达到最大重试次数,从日志中可看到首次尝试失败、第二次成功连接,最终服务正常启动。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 19:16:00