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

Docker Compose环境中Spark连接Kafka失败问题求助

问题:Spark无法连接Docker Compose部署的Kafka服务

执行Spark消费Kafka数据的作业时,出现以下错误:

24/12/13 23:28:02 WARN 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.8'

services:
  zookeeper:
    image: confluentinc/cp-zookeeper:latest
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181
      ZOOKEEPER_TICK_TIME: 2000

  kafka:
    image: confluentinc/cp-kafka:latest
    ports:
      - "9092:9092"
    environment:
      KAFKA_BROKER_ID: 1
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
      KAFKA_LISTENERS: PLAINTEXT://0.0.0.0:9092

  spark:
    image: bitnami/spark:latest
    ports:
      - "8080:8080"
    volumes:
      - 'C:/docker-proyecto/spark-data:/opt/spark'
    environment:
      - SPARK_MODE=master
      - SPARK_RPC_AUTHENTICATION_ENABLED=no

  spark-worker:
    image: bitnami/spark:latest
    environment:
      - SPARK_MODE=worker
      - SPARK_MASTER_URL=spark://spark-master:7077
      - SPARK_WORKER_MEMORY=1G
      - SPARK_WORKER_CORES=1
    depends_on:
      - spark
    volumes:
      - 'C:/docker-proyecto/spark-data:/opt/spark'

执行的spark-submit命令

spark-submit --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.3 /opt/spark/script.py

Spark脚本内容

kafka_servers = "localhost:9092"
stream_df = spark.readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", kafka_servers) \
    .option("subscribe", "topic1,topic2") \
    .option("startingOffsets", "earliest") \
    .load()

预期目标:Spark正常连接Kafka读取Topic数据,处理后存入HDFS。


问题原因与解决方法

1. 容器间网络访问错误

Spark容器内的localhost指向容器自身,而非宿主机,所以用localhost:9092无法访问到Kafka容器。Docker Compose会自动为服务创建内部DNS,直接使用Kafka服务名kafka即可访问对应容器。

2. Kafka监听地址配置缺陷

当前KAFKA_ADVERTISED_LISTENERS仅配置了宿主机可访问的地址,当Spark容器从Kafka获取broker地址时,得到的是localhost:9092,导致Spark容器尝试连接自己的9092端口,必然失败。需要为Kafka配置两个监听地址:一个供宿主机外部访问,一个供Docker内部容器间访问。


修改步骤

步骤1:更新docker-compose.yml的Kafka配置

修改Kafka服务的环境变量,添加内部监听地址:

kafka:
  image: confluentinc/cp-kafka:latest
  ports:
    - "9092:9092"
    - "9093:9093"  # 新增端口映射供内部访问
  environment:
    KAFKA_BROKER_ID: 1
    KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
    # 配置两个advertised listeners:外部用localhost,内部用服务名
    KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092,PLAINTEXT_INTERNAL://kafka:9093
    KAFKA_LISTENERS: PLAINTEXT://0.0.0.0:9092,PLAINTEXT_INTERNAL://0.0.0.0:9093
    KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_INTERNAL:PLAINTEXT
    KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT_INTERNAL

步骤2:修改Spark脚本的Kafka地址

将脚本中的kafka_servers改为Docker内部可访问的地址:

kafka_servers = "kafka:9093"  # 使用Kafka服务名+内部监听端口
stream_df = spark.readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", kafka_servers) \
    .option("subscribe", "topic1,topic2") \
    .option("startingOffsets", "earliest") \
    .load()

步骤3:验证连接

  • 重启Docker Compose服务:docker-compose down && docker-compose up -d
  • 进入Spark容器测试网络连通性:docker exec -it <spark-container-id> ping kafka,能ping通则内部网络正常
  • 可在Kafka容器内创建测试Topic:docker exec -it <kafka-container-id> kafka-topics.sh --create --topic topic1 --bootstrap-server localhost:9092 --partitions 1 --replication-factor 1

这样修改后,Spark容器就能通过内部网络正常连接到Kafka服务,同时宿主机依然可以通过localhost:9092访问Kafka。

内容的提问来源于stack exchange,提问作者Leynder Sánchez

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 15:13:10