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

Airflow中编程重跑失败DagRun无效问题求助

解决Airflow编程重跑失败DagRun的问题

原代码无法修改状态的核心原因是:仅调用set_state只是修改了内存中的对象状态,没有将变更提交到Airflow的元数据库。以下是修正后的实现方案:

修正思路

  • 获取目标DagRun对象后修改状态
  • 通过Airflow的SQLAlchemy会话提交变更
  • (可选)重置关联任务实例状态并触发完整重跑流程(仅改状态可能不足以触发任务执行)

完整代码示例

from airflow.models import DagRun, DagRunState
from airflow.utils.session import create_session

# 假设env_dag_id和dag_run_id已提前定义
with create_session() as session:
    # 从数据库查询目标DagRun
    dag_run = session.query(DagRun).filter(
        DagRun.dag_id == env_dag_id,
        DagRun.run_id == dag_run_id
    ).first()
    
    if dag_run and dag_run.state == DagRunState.FAILED:
        # 修改状态并绑定会话确保持久化
        dag_run.set_state(DagRunState.RUNNING, session=session)
        print(f"DagRun {dag_run_id} 状态已更新为RUNNING")

额外说明

如果需要真正触发任务重新执行,仅修改DagRun状态可能不够,还需重置关联TaskInstance的状态并触发运行:

from airflow.models import TaskInstance

# 重置该DagRun下所有失败/上游失败的任务实例
task_instances = session.query(TaskInstance).filter(
    TaskInstance.dag_id == env_dag_id,
    TaskInstance.run_id == dag_run_id,
    TaskInstance.state.in_([DagRunState.FAILED, DagRunState.UPSTREAM_FAILED])
).all()

for ti in task_instances:
    ti.set_state(None, session=session)

# 触发DAG运行(适配不同Airflow版本)
dag_run.trigger_run()
  • 使用create_session()上下文管理器可自动处理会话的提交与关闭,避免数据库连接泄漏。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 00:45:32