在MWAA Airflow中运行Spark遇磁盘空间不足问题求优化方案
问题描述
- 运行环境:亚马逊MWAA,环境类为mw1.large
- 配置详情:
- 最多1000个DAG,默认20并发任务
- 组件规格:2台Web服务器(2vCPU/4GB)、工作节点(4vCPU/8GB)、调度器(4vCPU/8GB)、数据库(2vCPU/8GB)
- 任务内容:读取S3中多个Parquet文件(总数据量约400GB),完成去重聚合后写入S3
- 错误信息:
INFO - Error processing PMC: An error occurred while calling o90.parquet.: org.apache.spark.SparkException: Job aborted due to stage failure: Task 1254 in stage 1.0 failed 1 times, most recent failure: Lost task 1254.0 in stage 1.0 (TID 1255) (executor driver): java.io.IOException: No space left on device
中文翻译:
INFO - 处理PMC时出错:调用o90.parquet时发生错误。: org.apache.spark.SparkException: 由于阶段失败导致作业中止:阶段1.0中的任务1254失败1次,最近一次失败:丢失阶段1.0中的任务1254.0(TID 1255)(执行器driver):java.io.IOException: 设备上没有剩余空间
优化建议
1. 迁移Spark临时存储到S3
MWAA工作节点本地磁盘空间有限,Spark处理大文件时的shuffle、临时数据会占满本地磁盘。直接将临时存储路径指向S3,彻底避开本地磁盘限制:
spark_conf = { "spark.hadoop.fs.s3a.buffer.dir": "s3://your-temp-bucket/spark-shuffle-tmp", "spark.local.dir": "s3://your-temp-bucket/spark-local-tmp", "spark.shuffle.service.enabled": "true", "spark.hadoop.fs.s3a.fast.upload": "true" }
确保MWAA执行角色拥有该临时桶的读写权限,任务结束后可自动或手动清理临时文件。
2. 优化数据处理逻辑,减少计算量
- 按需读取字段:仅加载去重聚合必需的字段,避免全量读取Parquet文件:
df = spark.read.parquet("s3://input-bucket/*").select("id", "value", "create_time") - 前置数据过滤:在去重前先过滤无效、过期数据,缩小处理数据集:
df = df.filter(df.create_time >= "2024-01-01" & df.is_valid == 1) - 精准去重:基于业务主键去重,避免全字段比对的资源浪费:
df = df.dropDuplicates(["id"])
3. 调整MWAA工作节点规格与并发
- 升级工作节点环境类:mw1.large默认本地磁盘约40GB,升级到mw1.xlarge(80GB本地盘)或更高规格,直接提升本地存储容量。
- 降低任务并发:临时将该DAG的并发数调低(如设为10),避免多个Spark任务同时占用磁盘资源,减少竞争。可通过DAG参数设置:
default_args = { # ...其他参数 "concurrency": 10 } dag = DAG("data_processing_dag", default_args=default_args, ...)
4. 优化Spark executor与分区配置
- 调整executor资源分配:给executor分配更多内存,减少磁盘spill的概率:
spark_conf = { "spark.executor.memory": "6g", "spark.executor.cores": "3", "spark.driver.memory": "6g", "spark.sql.shuffle.partitions": "200" # 减少shuffle分区数,降低临时文件总大小 } - 增大Parquet文件读取分区:如果输入Parquet文件过小,合并为大文件后再处理,减少任务数和临时文件数量。
5. 添加本地临时文件清理步骤
在DAG末尾添加清理任务,删除Spark生成的本地临时文件,避免残留文件占用磁盘:
from airflow.operators.bash import BashOperator cleanup_task = BashOperator( task_id="cleanup_spark_temp", bash_command="rm -rf /tmp/spark-* /tmp/hive-*", trigger_rule="all_done" )
内容的提问来源于stack exchange,提问作者Sourabh Somvanshi

