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

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 }}',
    }
)

关键修改点

  1. 优化查询逻辑:last_exec函数直接通过order_by获取最新的TaskInstance,提升查询效率。
  2. 完整更新TaskInstance状态:除了设置state为SUCCESS,还补充更新end_date和duration字段,确保状态数据完整。
  3. 触发调度器接管:找到对应DagRun并将其状态设为RUNNING,Airflow调度器会自动检查该DagRun的任务依赖,触发所有满足条件的下游任务(如示例中的run_them)。
  4. 增加异常处理:添加事务回滚逻辑,避免状态更新出现脏数据。

内容的提问来源于stack exchange,提问作者Nicolas Menettrier

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 18:50:26