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

Airflow:对应传感器失败时如何重试关联任务?

问题描述

我有一个Airflow DAG,包含两个核心任务:

  • start_job:负责通过REST API启动外部作业
  • wait_for_job:传感器任务,持续检查外部作业是否完成

需求是:当wait_for_job在配置的超时时间内未检测到作业完成(任务进入FAILED状态)时,同时重试start_job和wait_for_job,重新启动并等待新的作业实例。

当前遇到的问题:
通过wait_for_job的on_failure_callback将start_job的状态设为UP_FOR_RETRY后,start_job确实会重试执行,但重试完成后,下游的wait_for_job并未被触发。start_job的最后一条日志显示INFO - 0 downstream tasks scheduled from follow-on schedule check,而第一次执行start_job时能正常触发wait_for_job。

以下是去除REST API逻辑后的最小化DAG示例:

import time
import logging
from datetime import timedelta
from typing import Any, Dict, List, Optional

import pendulum
from sqlalchemy.orm.session import Session

from airflow.decorators import dag, task
from airflow.sensors.base import PokeReturnValue
from airflow.models import taskinstance
from airflow.utils.state import State
from airflow.utils.db import provide_session
from airflow.utils.session import NEW_SESSION, provide_session

logger = logging.getLogger("airflow.task")


@provide_session
def on_failure_callback(context: Dict[str, Any], session: Session = NEW_SESSION) -> None:
    logger.info(f"on_failure_callback()")

    start_job_task = _get_task(context, "start_job")
    wait_for_job_task = _get_task(context, "wait_for_job")

    # _clear_task(wait_for_job_task, session, context)
    # _clear_task(start_job_task, session, context)

    # start_job_task.set_state(State.FAILED)
    # wait_for_job_task.set_state(State.UP_FOR_RETRY)

    logging.info("set state of start_job_task to UP_FOR_RETRY ...")
    start_job_task.set_state(State.UP_FOR_RETRY)


def _clear_task(task, session, context):
    logger.info(f"run clear_task_instances() for task: {task.task_id}")
    taskinstance.clear_task_instances(
        tis=[task, ],
        session=session,
        dag=context["dag"])


def _get_task(
    context: Dict[str, Any],
    task_id: str,
    ) -> taskinstance.TaskInstance:
    
    task_instances: List[taskinstance.TaskInstance] = context["dag_run"].get_task_instances()
    logger.info(f"task_instances: {task_instances}")
    for ti in task_instances:
        logger.info(f"    ti.task_id: {ti.task_id}")
        if ti.task_id == task_id:
            return ti


@provide_session
def on_retry_callback(context: Dict[str, Any], session: Session = NEW_SESSION) -> None:
    print("on_retry_callback()")


@dag(
    schedule=None,
    start_date=pendulum.datetime(2023, 1, 1, tz="UTC"),
    catchup=False,
    tags=["job"],
)
def start_jobs_minimal_dag():

    @task(
        execution_timeout=timedelta(seconds=30),
        retries=3,
        retry_delay=timedelta(seconds=10),
    )
    def start_job():
        job_id = 1
        time.sleep(1)        
        return job_id
        
        
    @task.sensor(
        execution_timeout=timedelta(seconds=5),
        timeout=60,
        retries=3,
        retry_delay=timedelta(seconds=2),
        # mode='reschedule',
        on_failure_callback=on_failure_callback,
        on_retry_callback=on_retry_callback,
    )
    def wait_for_job(job_id: int) -> PokeReturnValue:
        logger.info(f"wait_for_job(): job_id: {job_id}")        
        time.sleep(2)        
        # 模拟传感器失败
        return PokeReturnValue(is_done=False, xcom_value=None)

    job_id = start_job()
    wait_for_job(job_id)
    
start_jobs_minimal_dag()
问题原因
  1. wait_for_job失败后处于FAILED状态(终端状态),Airflow不会自动重新调度已进入终端状态的任务,即使上游任务重试成功。
  2. 仅修改start_job的状态为UP_FOR_RETRY,但wait_for_job的FAILED状态未被清理,Airflow认为该任务已经完成,不会触发重新调度。
  3. 直接调用set_state修改任务状态的方式,没有更新Airflow调度所需的元数据(如依赖关系、任务实例的执行上下文),导致调度逻辑无法识别需要重新触发下游任务。
解决方案

要实现两个任务同时重试,需要在on_failure_callback中同时清理两个任务的状态,并触发重新调度:

  1. 使用Airflow官方的clear_task_instances方法,清除start_job和wait_for_job的任务实例状态,移除它们的终端状态标记。
  2. 将start_job设为UP_FOR_RETRY,或者直接让两个任务重新进入调度队列,确保Airflow能识别到依赖关系需要重新执行。

修改后的关键代码如下:

@provide_session
def on_failure_callback(context: Dict[str, Any], session: Session = NEW_SESSION) -> None:
    logger.info(f"on_failure_callback()")

    start_job_ti = _get_task(context, "start_job")
    wait_for_job_ti = _get_task(context, "wait_for_job")

    # 清除两个任务的状态,重置为未执行状态
    _clear_task(start_job_ti, session, context)
    _clear_task(wait_for_job_ti, session, context)

    # 将start_job设为UP_FOR_RETRY,触发重试
    start_job_ti.set_state(State.UP_FOR_RETRY, session=session)
    logger.info("Set start_job to UP_FOR_RETRY, wait_for_job will be scheduled after start_job succeeds")


def _clear_task(task_ti, session, context):
    logger.info(f"Clear task instance for: {task_ti.task_id}")
    # 使用clear_task_instances清除任务状态,包括XCom、日志等关联数据
    taskinstance.clear_task_instances(
        tis=[task_ti],
        session=session,
        dag=context["dag"],
        # 清除所有相关状态,包括失败标记
        include_upstream=False,
        include_downstream=False,
        reset_dag_runs=False
    )
关键说明
  • clear_task_instances会彻底清除任务实例的执行记录(包括状态、XCom、日志),确保任务可以被重新调度。
  • 先清除状态再设置UP_FOR_RETRY,可以保证start_job重试完成后,Airflow能正常识别到下游的wait_for_job需要重新执行。
  • 不需要手动设置wait_for_job的状态,因为start_job重试成功后,Airflow会根据依赖关系自动调度下游的wait_for_job。

内容的提问来源于stack exchange,提问作者dds-work

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 00:54:55