如何实现触发DAG标记主DAG上次运行为成功并每12小时重跑,单实例运行
解决主DAG多实例问题:标记旧实例为成功并保留最新运行
核心逻辑
在触发新的主DAG实例前,先查询主DAG的所有非成功状态运行实例,保留最新的一个,将其余旧实例的DAG Run及下属所有Task Instance标记为成功,同时暂停旧实例。
具体实现代码
在你的触发DAG中添加一个PythonOperator处理旧实例的状态更新,再串联触发主DAG的任务:
from airflow.models import DagRun, TaskInstance from airflow.utils.state import State from airflow.utils.db import create_session from airflow.operators.python import PythonOperator from airflow.operators.trigger_dagrun import TriggerDagRunOperator from airflow import DAG from datetime import datetime import datetime default_args = { 'owner': 'airflow', 'depends_on_past': False, 'start_date': datetime(2024, 1, 1), 'retries': 0, } def mark_old_dag_runs_success(dag_id): with create_session() as session: # 筛选主DAG的所有非成功运行实例,按执行时间倒序排列 dag_runs = session.query(DagRun)\ .filter(DagRun.dag_id == dag_id)\ .filter(DagRun.state != State.SUCCESS)\ .order_by(DagRun.execution_date.desc())\ .all() if len(dag_runs) <= 1: return # 处理除最新实例外的所有旧实例 for run in dag_runs[1:]: # 更新DAG Run状态为成功、设置结束时间并暂停 run.state = State.SUCCESS run.end_date = datetime.datetime.now() run.paused = True session.merge(run) # 更新该实例下所有任务的状态为成功 task_instances = session.query(TaskInstance)\ .filter(TaskInstance.dag_id == dag_id)\ .filter(TaskInstance.run_id == run.run_id)\ .all() for ti in task_instances: ti.state = State.SUCCESS ti.end_date = datetime.datetime.now() session.merge(ti) session.commit() # 定义触发DAG trigger_dag = DAG( 'trigger_main_dag', schedule_interval='@every 12h', catchup=False, default_args=default_args ) # 标记旧主DAG实例为成功的任务 mark_old_runs_task = PythonOperator( task_id='mark_old_main_dag_runs_success', python_callable=mark_old_dag_runs_success, op_kwargs={'dag_id': 'your_main_dag_id'}, # 替换为你的主DAG ID dag=trigger_dag ) # 触发新主DAG实例的任务 trigger_main_dag_task = TriggerDagRunOperator( task_id='trigger_main_dag_run', trigger_dag_id='your_main_dag_id', # 替换为你的主DAG ID dag=trigger_dag ) # 设置任务依赖:先处理旧实例,再触发新实例 mark_old_runs_task >> trigger_main_dag_task
关键注意事项
- 权限配置:确保触发DAG的执行账户拥有修改Airflow元数据的权限(如Admin角色),否则无法更新DAG Run和Task Instance的状态。
- 数据一致性:如果主DAG涉及数据写入等操作,强制标记旧实例为成功前需确认这些任务不会造成数据冲突,建议在测试环境验证后再部署到生产。
- 关闭自动补跑:主DAG和触发DAG都要设置
catchup=False,避免历史未执行的任务自动触发,导致额外实例。 - 状态筛选:代码中只筛选非成功状态的实例,已成功的实例不会被修改,避免干扰正常历史记录。
内容的提问来源于stack exchange,提问作者Ugur Selim Ozen
相关产品推荐
相关产品推荐

