如何识别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
相关产品推荐
相关产品推荐

