Airflow中如何将历史TaskInstance设为成功并续跑流水线?
问题:修改指定TaskInstance状态后自动续跑Airflow流水线
我想要创建一个可通过dag_id/task_id触发的DAG,目标是将指定DAG中最后执行的目标TaskInstance状态设为“success”,并从该节点开始自动续跑流水线。例如在示例流水线中,将“run_that”设为成功后自动运行“run_them”。
以下是我编写的代码:
import airflow from airflow.models import DagRun, TaskInstance, DagBag from airflow.operators.dagrun_operator import TriggerDagRunOperator from airflow.utils.trigger_rule import TriggerRule from airflow.utils.state import State from airflow.operators.python_operator import PythonOperator import pendulum from wrapper import ScriptLauncher, handleErrorSlack, handleErrorMail from datetime import timedelta, datetime default_args = { 'owner': 'tozzi', 'depends_on_past': False, 'start_date': pendulum.datetime(2022, 12, 19, tz='Europe/Paris'), 'retries': 0, 'retry_delay': timedelta(minutes=5), 'xcom_push': True, 'catchup': False, 'params': { 'dag_id': 'my_dag', 'task_id': 'run_that', } } def last_exec(dag_id, task_id, session): task_instances = ( session.query(TaskInstance) .filter(TaskInstance.dag_id == dag_id, TaskInstance.task_id == task_id) .all() ) task_instances.sort(key=lambda x: x.execution_date, reverse=True) if task_instances: return task_instances[0] return None def set_last_task_success(**kwargs): dag_id = kwargs['dag_id'] task_id = kwargs['task_id'] session = airflow.settings.Session() task_instance = last_exec(dag_id, task_id, session) if (task_instance is not None): task_instance.state = 'success' # task_instance = TaskInstance(task_id=task_id, execution_date=last_task_instance.execution_date) task_instance.run(session=session, ignore_ti_state=True, ignore_task_deps=True) session.commit() session.close() doc_md=f"""## Set the given task_id to success of the given dag_id""" # launched remotely launcher = ScriptLauncher(default_args, "@once", 'set_task_to_success', ['airflow'], doc_md) dag = launcher.dag; set_to_success = PythonOperator( task_id='set_to_success', provide_context=True, python_callable=set_last_task_success, dag=dag, op_kwargs={ 'dag_id': '{{ params.dag_id }}', 'task_id': '{{ params.task_id }}', } )
目前状态修改成功,但调用task_instance.run(...)时出现错误:AttributeError: 'TaskInstance' object has no attribute 'task',请问该如何修改才能在修改“run_that”状态后自动运行“run_them”任务?
解决方案
错误原因
从数据库直接查询得到的TaskInstance对象仅包含持久化的字段数据,并未关联到内存中对应的Task(DAG任务定义),而run()方法需要依赖task属性执行任务逻辑,因此触发报错。更符合Airflow设计的方式是通过更新状态后,让调度器自动接管后续任务的触发,而非手动调用run()。
修改后的代码
import airflow from airflow.models import DagRun, TaskInstance, DagBag from airflow.utils.state import State from airflow.operators.python_operator import PythonOperator import pendulum from wrapper import ScriptLauncher, handleErrorSlack, handleErrorMail from datetime import timedelta, datetime default_args = { 'owner': 'tozzi', 'depends_on_past': False, 'start_date': pendulum.datetime(2022, 12, 19, tz='Europe/Paris'), 'retries': 0, 'retry_delay': timedelta(minutes=5), 'xcom_push': True, 'catchup': False, 'params': { 'dag_id': 'my_dag', 'task_id': 'run_that', } } def last_exec(dag_id, task_id, session): # 优化查询:直接通过order_by获取最新的TaskInstance,避免全量查询后排序 return ( session.query(TaskInstance) .filter(TaskInstance.dag_id == dag_id, TaskInstance.task_id == task_id) .order_by(TaskInstance.execution_date.desc()) .first() ) def set_last_task_success(**kwargs): dag_id = kwargs['dag_id'] task_id = kwargs['task_id'] tz = pendulum.timezone('Europe/Paris') session = airflow.settings.Session() try: # 获取目标TaskInstance task_instance = last_exec(dag_id, task_id, session) if not task_instance: print(f"No TaskInstance found for dag_id={dag_id}, task_id={task_id}") return # 加载目标DAG验证任务存在 dag_bag = DagBag() target_dag = dag_bag.get_dag(dag_id) if not target_dag: raise ValueError(f"DAG {dag_id} not found in DagBag") if task_id not in target_dag.task_ids: raise ValueError(f"Task {task_id} not found in DAG {dag_id}") # 更新TaskInstance状态及相关字段 task_instance.state = State.SUCCESS task_instance.end_date = datetime.now(tz) if task_instance.start_date: task_instance.duration = task_instance.end_date - task_instance.start_date session.commit() # 找到对应的DagRun,标记为RUNNING以触发调度器重新处理 dag_run = ( session.query(DagRun) .filter( DagRun.dag_id == dag_id, DagRun.execution_date == task_instance.execution_date ) .first() ) if dag_run: dag_run.state = State.RUNNING session.commit() print(f"Updated DagRun {dag_run} to RUNNING, scheduler will trigger downstream tasks") except Exception as e: session.rollback() raise e finally: session.close() doc_md=f"""## Set the given task_id to success of the given dag_id""" # launched remotely launcher = ScriptLauncher(default_args, "@once", 'set_task_to_success', ['airflow'], doc_md) dag = launcher.dag; set_to_success = PythonOperator( task_id='set_to_success', provide_context=True, python_callable=set_last_task_success, dag=dag, op_kwargs={ 'dag_id': '{{ params.dag_id }}', 'task_id': '{{ params.task_id }}', } )
关键修改点
- 优化查询逻辑:
last_exec函数直接通过order_by获取最新的TaskInstance,提升查询效率。 - 完整更新TaskInstance状态:除了设置
state为SUCCESS,还补充更新end_date和duration字段,确保状态数据完整。 - 触发调度器接管:找到对应DagRun并将其状态设为RUNNING,Airflow调度器会自动检查该DagRun的任务依赖,触发所有满足条件的下游任务(如示例中的
run_them)。 - 增加异常处理:添加事务回滚逻辑,避免状态更新出现脏数据。
内容的提问来源于stack exchange,提问作者Nicolas Menettrier
相关产品推荐
相关产品推荐

