【问题标题】:How to check different run times of a task in a DAG in an External Sensor Airflow如何在外部传感器气流中检查 DAG 中任务的不同运行时间
【发布时间】:2021-11-03 13:12:33
【问题描述】:

假设我有一个 DAG A 包括一些任务,而这些任务依赖于 DAG B 中另一个任务上的某个外部传感器。

例如我想在10:00检查DAG B中的一个任务的状态,如果这次运行成功,那么中的任务DAG A 可以运行。 但是现在由于某种原因,10:00DAG B中的任务失败了,但是11:00运行了相同的任务> 成功了。

问题是 DAG A 中的任务将永远挂起,因为 DAG B 中的任务在 10:00 失败。但是如果下次运行成功就可以了。

如何在外部传感器气流中实现这样的功能,检查另一个 DAG 中下一次运行时的状态,如果成功,那么我的任务可以毫无问题地运行?

P.S:由于某些原因我无法使用重试!

提前谢谢你。

【问题讨论】:

    标签: python python-3.x airflow airflow-scheduler


    【解决方案1】:

    我自己找到了解决方案。这可能不是最好的方法,但它确实有效。 在这种情况下,我们可以为 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)
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-12-31
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多