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

如何识别Airflow中所有损坏的DAG并删除对应存储桶的Python文件

GCP Airflow 自动清理损坏DAG实现方案

核心逻辑

  • 从Airflow元数据库的import_error表拉取所有存在导入错误的DAG对应的本地文件路径
  • 路径转换为GCS DAG存储桶对应的对象路径
  • 验证文件状态(排除近期新上传还未完成解析的文件)后删除GCS对应文件
  • 可选清理Airflow元数据库中残留的损坏DAG记录

前置依赖

  • Airflow Worker服务账号已授予目标DAG存储桶的storage.objects.delete权限
  • 已安装google-cloud-storage、apache-airflow-providers-google依赖包(Cloud Composer默认预装)
  • 适配Airflow 2.x版本,1.x版本需调整元数据查询字段

完整DAG代码

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.utils.session import provide_session
from airflow.models import ImportError, DagModel
from google.cloud import storage
from datetime import datetime, timedelta
import os

# 自定义配置项
DAG_BUCKET = "你的DAG存储桶名称"  # 示例:us-central1-xxx-composer-bucket
# Cloud Composer中DAG本地路径前缀为/home/airflow/gcs/dags/,自行替换为你环境的对应前缀
LOCAL_DAG_PREFIX = "/home/airflow/gcs/dags/"
# 保护期:最近N小时内上传的文件不删除,避免误删未完成解析的正常DAG
PROTECT_HOURS = 1
# 干运行模式:开启时只打印要删除的文件,不执行真实删除操作,测试阶段建议开启
DRY_RUN = True
# 排除路径列表:不需要清理的DAG路径,比如动态DAG生成脚本的路径
EXCLUDE_PATHS = ["dynamic_dags/"]

default_args = {
    "owner": "airflow",
    "retries": 1,
    "retry_delay": timedelta(minutes=5),
    "start_date": datetime(2024, 1, 1),
}

@provide_session
def get_broken_dag_files(session=None, **context):
    # 查询所有有导入错误的DAG文件路径
    broken_files = session.query(ImportError.filename).distinct().all()
    # 过滤掉已删除的文件和排除路径的文件
    valid_broken_files = []
    for f in broken_files:
        file_path = f[0]
        # 过滤排除路径
        if any([file_path.startswith(os.path.join(LOCAL_DAG_PREFIX, p)) for p in EXCLUDE_PATHS]):
            continue
        # 转换为GCS对象路径
        if file_path.startswith(LOCAL_DAG_PREFIX):
            gcs_path = file_path.replace(LOCAL_DAG_PREFIX, "")
            valid_broken_files.append(gcs_path)
    # 去重后返回
    return list(set(valid_broken_files))

def delete_gcs_files(gcs_files, **context):
    if not gcs_files:
        print("未检测到损坏的DAG文件,无需清理")
        return
    storage_client = storage.Client()
    bucket = storage_client.bucket(DAG_BUCKET)
    protect_time = datetime.utcnow() - timedelta(hours=PROTECT_HOURS)
    deleted_count = 0
    for gcs_path in gcs_files:
        blob = bucket.blob(gcs_path)
        if not blob.exists():
            print(f"文件{gcs_path}不存在,跳过")
            continue
        # 检查文件更新时间,在保护期内的跳过
        if blob.updated.replace(tzinfo=None) > protect_time:
            print(f"文件{gcs_path}在保护期内,跳过")
            continue
        if DRY_RUN:
            print(f"干运行模式:待删除文件 {gcs_path}")
            deleted_count +=1
            continue
        # 执行删除
        blob.delete()
        print(f"已删除损坏DAG文件:{gcs_path}")
        deleted_count +=1
    print(f"本次共处理{deleted_count}个损坏DAG文件")

@provide_session
def clear_broken_dag_records(session=None, **context):
    # 清理元数据中已不存在的DAG记录,可选执行
    if DRY_RUN:
        print("干运行模式:跳过元数据清理")
        return
    session.query(DagModel).filter(DagModel.is_active == False, DagModel.has_import_errors == True).delete()
    session.commit()
    print("已清理元数据库中损坏DAG残留记录")

with DAG(
    "auto_clean_broken_dags",
    default_args=default_args,
    schedule_interval="0 2 * * *", # 每天凌晨2点运行,可自行调整
    catchup=False,
    tags=["admin", "maintenance"],
) as dag:
    get_broken_files_task = PythonOperator(
        task_id="get_broken_dag_files",
        python_callable=get_broken_dag_files,
    )
    delete_files_task = PythonOperator(
        task_id="delete_gcs_files",
        python_callable=delete_gcs_files,
        op_kwargs={"gcs_files": "{{ ti.xcom_pull(task_ids='get_broken_dag_files') }}"},
    )
    clear_records_task = PythonOperator(
        task_id="clear_broken_dag_records",
        python_callable=clear_broken_dag_records,
    )
    get_broken_files_task >> delete_files_task >> clear_records_task

使用注意事项

  • 首次部署后务必保持DRY_RUN = True运行至少1次,确认待删除的文件列表符合预期后再关闭干运行模式,避免误删正常DAG
  • 保护期PROTECT_HOURS建议至少设置为1小时,避免刚上传还未完成全量解析的正常DAG被误识别为损坏DAG
  • 若存在动态生成DAG的场景,请将动态DAG的生成脚本路径或者输出路径加入EXCLUDE_PATHS列表
  • 脚本默认每天凌晨2点运行,可根据自身需求调整schedule_interval参数

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 11:06:04