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

AWS托管Apache Airflow因磁盘耗尽无法运行,寻求恢复方案

解决AWS托管Airflow磁盘耗尽导致DAG失败的问题

一、紧急恢复步骤

1. 自行部署EC2场景

首先通过SSH或Session Manager登录Airflow所在EC2实例,执行以下操作:

  • 查看磁盘占用情况,定位满额挂载点:
    df -h
    
  • 找出占用空间最大的目录:
    du -h --max-depth=1 / | sort -hr
    
  • 优先清理Airflow核心占用文件:
    • 删除7天前的旧日志(默认日志路径为$AIRFLOW_HOME/logs):
      find $AIRFLOW_HOME/logs -type f -mtime +7 -delete
      
    • 紧急清理临时日志(若日志量过大影响运行):
      rm -rf $AIRFLOW_HOME/logs/*/*/*
      
    • 清理系统临时文件:
      find /tmp -type f -mtime +1 -delete
      rm -rf $AIRFLOW_HOME/tmp/*
      
    • 删除DAG任务未清理的下载文件(替换为你的实际下载路径):
      rm -rf /path/to/your/downloaded/files/*
      

2. MWAA托管场景

由于无法直接登录实例,按以下步骤操作:

  • 暂停所有DAG,避免新任务持续占用磁盘
  • 若使用Fargate Worker:等待MWAA自动回收临时Worker实例,销毁后磁盘空间会自动释放
  • 若使用EC2 Worker:在MWAA控制台调整Worker自动伸缩配置,先缩容至0再扩容,新启动的Worker实例磁盘为干净状态

二、彻底避免重复问题的方案

1. 完善DAG的强制清理逻辑

确保任务无论成功/失败都清理本地文件,示例代码:

from airflow.decorators import task
import os
import shutil
from airflow.models import Variable

@task
def download_upload_task():
    download_path = "/tmp/airflow_task_data"
    try:
        # 数据下载逻辑
        # S3上传逻辑
        pass
    finally:
        # 无论执行结果如何都清理本地文件
        if os.path.exists(download_path):
            shutil.rmtree(download_path)

或通过回调函数统一处理:

def cleanup_files(**context):
    download_path = context["ti"].xcom_pull(key="download_path")
    if download_path and os.path.exists(download_path):
        shutil.rmtree(download_path)

with DAG(...) as dag:
    upload_task = PythonOperator(
        task_id="download_upload",
        python_callable=download_upload_task,
        on_success_callback=cleanup_files,
        on_failure_callback=cleanup_files
    )

2. 配置日志自动过期

  • 自行部署:修改airflow.cfg,设置日志保留天数:
    [logging]
    log_retention_days = 30
    
  • MWAA:在环境配置的airflow.cfg覆盖项中添加上述配置

3. 监控磁盘使用率

  • 自行部署EC2:在CloudWatch中创建磁盘使用率告警,阈值设为80%,触发后及时处理
  • MWAA:监控CloudWatch指标WorkerDiskUtilization,设置告警规则

4. 流式上传避免本地存储(大文件场景)

直接通过流式下载+上传到S3,无需落地本地文件:

import boto3
import requests

@task
def stream_upload_to_s3():
    data_url = Variable.get("data_source_url")
    s3_client = boto3.client("s3")
    with requests.get(data_url, stream=True) as resp:
        s3_client.upload_fileobj(resp.raw, "your-s3-bucket", "target/file/path")

内容的提问来源于stack exchange,提问作者Meio

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 17:40:59