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

在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 23:45:01