Docker中Spark作业输出后永久停滞问题排查求助
Spark作业执行完成后进程停滞的解决方法
问题背景
使用bitnami/spark:3镜像构建Docker容器,通过Docker Compose部署Spark Master和Worker节点。作业已输出预期结果,但执行后进程永久停滞,需手动终止,无法直接与Airflow集成。
相关配置与代码
Dockerfile
FROM bitnami/spark:3 USER root RUN curl https://repo1.maven.org/maven2/com/amazonaws/aws-java-sdk-bundle/1.12.231/aws-java-sdk-bundle-1.12.231.jar --output /opt/bitnami/spark/jars/aws-java-sdk-bundle-1.12.231.jar RUN curl https://repo1.maven.org/maven2/net/java/dev/jets3t/jets3t/0.9.4/jets3t-0.9.4.jar --output /opt/bitnami/spark/jars/jets3t-0.9.4.jar RUN curl https://s3.amazonaws.com/redshift-downloads/drivers/jdbc/2.1.0.10/redshift-jdbc42-2.1.0.10.jar --output /opt/bitnami/spark/jars/redshift-jdbc42-2.1.0.10.jar RUN curl https://storage.googleapis.com/spark-lib/bigquery/spark-bigquery-with-dependencies_2.12-0.28.0.jar --output /opt/bitnami/spark/jars/spark-bigquery-with-dependencies_2.12-0.28.0.jar COPY ./requirements_copy.txt / RUN pip install -r /requirements_copy.txt
docker-compose.yml
version: '3' services: spark: image: spark-air:latest environment: - SPARK_MODE=master - SPARK_RPC_AUTHENTICATION_ENABLED=no - SPARK_RPC_ENCRYPTION_ENABLED=no - SPARK_LOCAL_STORAGE_ENCRYPTION_ENABLED=no - SPARK_SSL_ENABLED=no - AWS_ACCESS_KEY=${AWS_ACCESS_KEY} - AWS_SECRET_KEY=${AWS_SECRET_KEY} volumes: - ./dags:/opt/bitnami/spark/dags/:rw ports: - '8090:8080' spark-worker: image: spark-air:latest environment: - SPARK_MODE=worker - SPARK_MASTER_URL=spark://spark:7077 - SPARK_WORKER_MEMORY=1G - SPARK_WORKER_CORES=1 - SPARK_RPC_AUTHENTICATION_ENABLED=no - SPARK_RPC_ENCRYPTION_ENABLED=no - SPARK_LOCAL_STORAGE_ENCRYPTION_ENABLED=no - SPARK_SSL_ENABLED=no - AWS_ACCESS_KEY=${AWS_ACCESS_KEY} - AWS_SECRET_KEY=${AWS_SECRET_KEY} volumes: - ./dags:/opt/bitnami/spark/dags/:rw
Spark Python代码(原代码)
from pyspark import SparkContext sc = SparkContext("local", "First App") data = [{"Category": 'A', "ID": 1, "Value": 121.44, "Truth": True}, {"Category": 'B', "ID": 2, "Value": 300.01, "Truth": False}, {"Category": 'C', "ID": 3, "Value": 10.99, "Truth": None}, {"Category": 'E', "ID": 4, "Value": 33.87, "Truth": True}] df = sc.parallelize(data) df = df.collect() print(df)
作业日志
23/02/15 05:32:58 INFO Executor: Finished task 0.0 in stage 0.0 (TID 0). 1090 bytes result sent to driver 23/02/15 05:32:58 DEBUG ExecutorMetricsPoller: stageTCMP: (0, 0) -> 0 23/02/15 05:32:58 DEBUG TaskSchedulerImpl: parentName: , name: TaskSet_0.0, runningTasks: 0 23/02/15 05:32:58 DEBUG TaskSetManager: No tasks for locality level NO_PREF, so moving to locality level ANY 23/02/15 05:32:58 INFO TaskSetManager: Finished task 0.0 in stage 0.0 (TID 0) in 1867 ms on d981929c6c2c (executor driver) (1/1) 23/02/15 05:32:58 INFO TaskSchedulerImpl: Removed TaskSet 0.0, whose tasks have all completed, from pool 23/02/15 05:32:58 INFO DAGScheduler: ResultStage 0 (collect at /opt/bitnami/spark/dags/test.py:12) finished in 3.275 s 23/02/15 05:32:58 DEBUG DAGScheduler: After removal of stage 0, remaining stages = 0 23/02/15 05:32:58 INFO DAGScheduler: Job 0 is finished. Cancelling potential speculative or zombie tasks for this job 23/02/15 05:32:58 INFO TaskSchedulerImpl: Killing all running tasks in stage 0: Stage finished 23/02/15 05:32:58 INFO DAGScheduler: Job 0 finished: collect at /opt/bitnami/spark/dags/test.py:12, took 3.794097 s [{'Category': 'A', 'ID': 1, 'Value': 121.44, 'Truth': True}, {'Category': 'B', 'ID': 2, 'Value': 300.01, 'Truth': False}, {'Category': 'C', 'ID': 3, 'Value': 10.99, 'Truth': None}, {'Category': 'E', 'ID': 4, 'Value': 33.87, 'Truth': True}] 23/02/15 05:33:04 DEBUG ExecutorMetricsPoller: removing (0, 0) from stageTCMP
解决方法
1. 显式停止Spark上下文
原代码仅创建SparkContext但未主动停止,导致进程无法自动退出。在代码末尾添加停止命令:
# 停止SparkContext sc.stop()
2. 改用SparkSession(推荐)
Spark 2.x及以上版本推荐使用SparkSession替代底层SparkContext,它会自动管理上下文生命周期,也便于扩展功能:
from pyspark.sql import SparkSession # 初始化SparkSession spark = SparkSession.builder.appName("First App").getOrCreate() sc = spark.sparkContext data = [{"Category": 'A', "ID": 1, "Value": 121.44, "Truth": True}, {"Category": 'B', "ID": 2, "Value": 300.01, "Truth": False}, {"Category": 'C', "ID": 3, "Value": 10.99, "Truth": None}, {"Category": 'E', "ID": 4, "Value": 33.87, "Truth": True}] df = sc.parallelize(data) df = df.collect() print(df) # 停止SparkSession(自动停止关联的SparkContext) spark.stop()
3. 调整作业提交方式
当前Docker Compose配置启动的是长期运行的Spark集群服务,但Airflow集成需要一次性作业执行后自动退出。可通过spark-submit直接提交作业到集群:
# 替换为你的Docker Compose网络名称 docker run --network <your-compose-network-name> spark-air:latest spark-submit --master spark://spark:7077 /opt/bitnami/spark/dags/test.py
作业执行完成后,spark-submit进程结束,容器会自动退出,无需手动终止。
4. 匹配集群运行模式
原代码使用local模式初始化SparkContext,但实际已部署集群,建议修改为连接到集群Master:
# 连接到Spark集群Master sc = SparkContext("spark://spark:7077", "First App")
或在SparkSession中配置:
spark = SparkSession.builder.appName("First App").master("spark://spark:7077").getOrCreate()
内容的提问来源于stack exchange,提问作者light yamagi
相关产品推荐
相关产品推荐

