我自己找到了解决方案。这可能不是最好的方法,但它确实有效。
在这种情况下,我们可以为 on_failure_callback 定义一个函数,并为我们的 ExternalSensor 设置一个 timeout,当达到超时时,我们检查下一次运行的任务另一个 DAG,如果成功,我们将 ExternalSensor 的状态设置为SUCCESS,以便依赖此传感器的其他任务可以正常运行。
这是这种方法的代码:
from airflow.utils.state import State
from airflow.sensors.external_task_sensor import ExternalTaskSensor
from airflow.exceptions import AirflowSensorTimeout
from datetime import datetime, timedelta, timezone
from airflow.api.common.experimental.get_task_instance import get_task_instance
from dateutil.parser import parse
from functools import partial
def _failure_callback(task_id, dag_id, execution_date, context):
if isinstance(context['exception'], AirflowSensorTimeout):
sensor_instance = context['task_instance']
next_execution_date = parse(context['ts']) + -(execution_date) + timedelta(hours=1)
ti = get_task_instance(dag_id=dag_id, task_id=task_id, execution_date=next_execution_date)
if ti.current_state() == 'success':
sensor_instance.set_state(State.SUCCESS)
sensor = ExternalTaskSensor(external_task_id='external_task_id',
task_id='sensor',
external_dag_id='external_dag_id',
execution_delta=timedelta(hours=-24) + timedelta(minutes=-30),
timeout=5,
on_failure_callback=partial(_failure_callback, 'external_task_id', 'external_dag_id', timedelta(hours=-24) + timedelta(minutes=-30)),
dag=dag)