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
相关产品推荐
相关产品推荐

