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

如何实现触发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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 05:25:17