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

Docker Compose部署Spark-Kafka集群消费任务卡住问题求助

Spark与Kafka Docker集群消费任务卡住问题解决

问题概述

通过Docker Compose搭建Spark与Kafka集群后,提交Kafka主题消费任务时任务始终卡住无进展,但该任务在本地独立Spark集群可正常运行。Spark镜像基于bitnami/spark:3.5构建。

Docker Compose配置

version: '3'
services:
  spark-master:
    image: spark-cluster:1.0
    container_name: spark-master
    hostname: spark-master
    environment:
      - SPARK_MODE=master
      - SPARK_MASTER_HOSTNAME=spark-master
      - SPARK_MASTER_PORT=7077
    ports:
      - "4040:4040"
      - "6066:6066"
      - "7077:7077"
      - "8080:8080"
    deploy:
      resources:
        limits:
          cpus: "0.5" # Adjust as needed
          memory: 256M  # Adjust as needed
    networks:
      - spark-network

  spark-worker-1:
    image: spark-cluster:1.0
    container_name: spark-worker-1
    environment:
      - SPARK_MODE=worker
      - SPARK_MASTER_URL=spark://spark-master:7077
      - SPARK_WORKER_MEMORY = 4g
    deploy:
      resources:
        limits:
          cpus: "4" # Adjust as needed
          memory: 5G  # Adjust as needed
    depends_on:
      - spark-master
    networks:
      - spark-network

  zookeeper:
    image: bitnami/zookeeper:3.9
    container_name: zookeeper-server
    restart: always
    ports:
      - "2181:2181"
    environment:
      - ALLOW_ANONYMOUS_LOGIN=yes

  kafka1:
    image: bitnami/kafka:3.5
    container_name: broker-1
    ports:
      - "9093:9093"
    environment:
      - KAFKA_BROKER_ID=1
      - KAFKA_ZOOKEEPER_CONNECT=zookeeper:2181
      - KAFKA_CFG_LISTENERS=PLAINTEXT://:9093,INTERNAL://:9092
      - KAFKA_CFG_ADVERTISED_LISTENERS=PLAINTEXT://localhost:9093,INTERNAL://:9092
      - KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP=PLAINTEXT:PLAINTEXT,INTERNAL:PLAINTEXT
      - KAFKA_CFG_INTER_BROKER_LISTENER_NAME=INTERNAL
    depends_on:
      - zookeeper 

  kafka2:
    image: bitnami/kafka:3.5
    container_name: broker-2
    ports:
      - "9094:9094"
    environment:
      - KAFKA_BROKER_ID=2
      - KAFKA_ZOOKEEPER_CONNECT=zookeeper:2181
      - KAFKA_CFG_LISTENERS=PLAINTEXT://:9094,INTERNAL://:9092
      - KAFKA_CFG_ADVERTISED_LISTENERS=PLAINTEXT://localhost:9094,INTERNAL://:9092
      - KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP=PLAINTEXT:PLAINTEXT,INTERNAL:PLAINTEXT
      - KAFKA_CFG_INTER_BROKER_LISTENER_NAME=INTERNAL
    depends_on:
      - zookeeper  

  kafka3:
    image: bitnami/kafka:3.5
    container_name: broker-3
    ports:
      - "9095:9095"
    environment:
      - KAFKA_BROKER_ID=3
      - KAFKA_ZOOKEEPER_CONNECT=zookeeper:2181
      - KAFKA_CFG_LISTENERS=PLAINTEXT://:9095,INTERNAL://:9092
      - KAFKA_CFG_ADVERTISED_LISTENERS=PLAINTEXT://localhost:9095,INTERNAL://:9092
      - KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP=PLAINTEXT:PLAINTEXT,INTERNAL:PLAINTEXT
      - KAFKA_CFG_INTER_BROKER_LISTENER_NAME=INTERNAL
    depends_on:
      - zookeeper  
  
networks:
  spark-network:
    driver: bridge

Consumer代码

from pyspark.sql import SparkSession
from pyspark.sql.functions import explode
from pyspark.sql.functions import split
import os
scala_version = '2.12'
spark_version = '3.5.0'
packages = [
    f'org.apache.spark:spark-sql-kafka-0-10_{scala_version}:{spark_version}',
    'org.apache.kafka:kafka-clients:3.5.0',
    'org.apache.hadoop:hadoop-client:3.0.0',
]

spark = SparkSession.builder \
   .master("spark://172.18.32.1:7077") \
   .appName("kafka-example") \
   .config("spark.jars.packages", ",".join(packages)) \
   .getOrCreate()

spark.sparkContext.setLogLevel("ERROR")
# Kafka broker地址
# kafka_brokers = "localhost:9092"
kafka_brokers = "localhost:9095"

# 定义要读取数据的Kafka主题
kafka_topic = "mytopic1"

# 从Kafka读取数据
df = spark \
  .readStream \
  .format("kafka") \
  .option("kafka.bootstrap.servers", kafka_brokers) \
  .option("subscribe", kafka_topic) \
  .option("startingOffsets", "earliest") \
  .load()

# 展示Kafka中的数据
castDf = df .selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
query = castDf.writeStream \
    .outputMode("append") \
    .format("console") \
    .start()

query.awaitTermination()

运行输出

Setting default log level to "WARN".
To adjust logging level use sc.setLogLevel(newLevel). For SparkR, use setLogLevel(newLevel).
23/11/28 23:08:17 WARN Utils: Service 'SparkUI' could not bind on port 4040. Attempting port 4041.

[Stage 0:>                                                          (0 + 0) / 1]

解决方案

1. 统一容器网络

当前Docker Compose中,Spark集群在spark-network网络,但Zookeeper和Kafka集群未加入该网络,导致Spark容器无法和Kafka容器通信。需为Zookeeper和所有Kafka节点添加networks配置:

# 在zookeeper服务中添加
networks:
  - spark-network

# 在kafka1、kafka2、kafka3服务中分别添加
networks:
  - spark-network

2. 修正Kafka监听器配置

Kafka的内部监听器INTERNAL的advertised.listeners配置为:9092,容器间无法通过该地址正确寻址。需修改为容器名+端口:

# 以kafka1为例,kafka2、kafka3做相同修改
environment:
  - KAFKA_CFG_ADVERTISED_LISTENERS=PLAINTEXT://localhost:9093,INTERNAL://broker-1:9092

3. 调整Spark代码中的Kafka地址

Spark容器在集群网络内,应使用Kafka的内部监听器地址(容器名+9092端口),而非宿主机的localhost地址:

kafka_brokers = "broker-1:9092,broker-2:9092,broker-3:9092"

4. 修正Spark Master地址

代码中使用固定IP172.18.32.1易导致网络问题,应使用Docker容器名spark-master:

spark = SparkSession.builder \
   .master("spark://spark-master:7077") \
   .appName("kafka-example") \
   .config("spark.jars.packages", ",".join(packages)) \
   .getOrCreate()

5. 修复Spark Worker环境变量格式

SPARK_WORKER_MEMORY后存在空格,导致变量无法正确解析,应删除空格:

environment:
  - SPARK_WORKER_MEMORY=4g

6. 提升Spark Master资源限制

当前Spark Master的CPU(0.5核)和内存(256M)限制过低,可能导致任务调度失败,建议调高:

deploy:
  resources:
    limits:
      cpus: "1"
      memory: 1G

内容的提问来源于stack exchange,提问作者Nguyễn Quốc Nhật Minh

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 02:49:52