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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 20:25:17