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
相关产品推荐
相关产品推荐

