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

