【问题标题】:Airflow: Master Dag with ExternalTaskSensor gets stuck forever气流:带有 ExternalTask​​Sensor 的 Master Dag 永远卡住了
【发布时间】:2020-06-25 20:32:37
【问题描述】:

要求是让 DAG 一个接一个地运行,并且每个 DAG 都成功

我有一个 Master DAG,我在其中调用所有 DAG 以依次执行

此外,在每个 dag_A、dag_B、dag_C 中,我必须给定 schedule_interval = None 并在 GUI 中手动打开

我正在使用 ExternalTask​​Sensor,因为即使在第一个 dag_A 中的所有任务完成之前,它也会启动第二个 dag_B,为避免此类问题,我正在使用 ExternalTask​​Sensor。如果有更好的实现,请告诉我

不知道我在这里缺少什么

代码:master_dag.py

import datetime
import os
from datetime import timedelta

from airflow.models import DAG, Variable

from airflow.operators.dagrun_operator import TriggerDagRunOperator
from airflow.operators.sensors import ExternalTaskSensor

default_args = {
        'owner': 'airflow',
        'start_date': datetime.datetime(2020, 1, 7),
        'provide_context': True,
        'execution_timeout': None,
        'retries': 0,
        'retry_delay': timedelta(minutes=3),
        'retry_exponential_backoff': True,
        'email_on_retry': False,
    }


dag = DAG(
        dag_id='master_dag',
        schedule_interval='7 3 * * *',
        default_args=default_args,
        max_active_runs=1,
        catchup=False,
    )

trigger_dag_A = TriggerDagRunOperator(
    task_id='trigger_dag_A',
        trigger_dag_id='dag_A',
        dag=dag,
    )

wait_for_dag_A = ExternalTaskSensor(
    task_id='wait_for_dag_A',
    external_dag_id='dag_A',
    external_task_id='proc_success',
    poke_interval=60,
    allowed_states=['success'],
    dag=dag,
    )

trigger_dag_B = TriggerDagRunOperator(
        task_id='trigger_dag_B',
        trigger_dag_id='dag_B',
        dag=dag,
    )

wait_for_dag_B = ExternalTaskSensor(
    task_id='wait_for_dag_B',
    external_dag_id='dag_B',
    external_task_id='proc_success',
    poke_interval=60,
    allowed_states=['success'],
    dag=dag)

trigger_dag_C = TriggerDagRunOperator(
        task_id='trigger_dag_C',
        trigger_dag_id='dag_C',
        dag=dag,
    )

trigger_dag_A >> wait_dag_A >> trigger_dag_B >> wait_dag_B >> trigger_dag_C

每个 DAG 都有多个任务在运行,最后一个任务是 proc_success

【问题讨论】:

    标签: airflow


    【解决方案1】:

    背景

    • ExternalTaskSensor 通过分别轮询外部DAGtaskDagRun / TaskInstance 的状态来工作(根据external_task_id 是否通过)
    • 现在,由于单个 DAG 可以有多个活动 DagRuns,因此必须告知传感器它应该感知哪些运行/实例
    • 为此,它使用execution_date 作为区分标准。这可以(仅)以其中一种以下两种方式表达
    :param execution_delta: time difference with the previous execution to
        look at, the default is the same execution_date as the current task or DAG.
        For yesterday, use [positive!] datetime.timedelta(days=1). Either
        execution_delta or execution_date_fn can be passed to
        ExternalTaskSensor, but not both.
    :type execution_delta: datetime.timedelta
    
    :param execution_date_fn: function that receives the current execution date
        and returns the desired execution dates to query. Either execution_delta
        or execution_date_fn can be passed to ExternalTaskSensor, but not both.
    :type execution_date_fn: callable
    

    你的实现中的问题

    • 在您的ExternalTaskSensors 中,您没有通过execution_date_fnexecution_delta 参数中的任何一个
    • 因此,传感器 picks up its own execution_date 轮询 child DAGs 的 DagRuns,从而卡住(显然,您的父/协调器 DAG 的 execution_date 将是不同于子 DAG)
    @provide_session
    def poke(self, context, session=None):
        if self.execution_delta:
            dttm = context['execution_date'] - self.execution_delta
        elif self.execution_date_fn:
            dttm = self.execution_date_fn(context['execution_date'])
        else:
            # if neither of above is passed, use current DAG's execution date
            dttm = context['execution_date']
    

    更多提示

    • 你可以跳过external_task_id;当你这样做时,ExternalTaskSensor 实际上变成了一个ExternalDagSensor。当您的子 DAG(A、B 和 C)有多个结束任务时,这尤其有用(因此完成任何一个结束任务并不能保证完成整个 DAG)
    • 也看看这个讨论:Wiring top-level DAGs together


    EDIT-1

    事后想来,我最初的判断似乎是错误的;特别是以下陈述不成立。

    显然,您的父/协调器 DAG 的 execution_date 将是 与子 DAG 不同

    查看the source,很明显TriggerDagRunOperator 将自己的execution_date 传递给子DagRun,这意味着ExternalTaskSensor 应该能够感知 DAG 或这是任务。

     trigger_dag(
                dag_id=self.trigger_dag_id,
                run_id=run_id,
                conf=self.conf,
                # own execution date passed to child DAG
                execution_date=self.execution_date,
                replace_microseconds=False,
            )
    

    所以这个解释不成立。

    我建议你去

    • 在 UI 中或通过查询元数据库来检查您触发的子 DAG 的 execution_date/您正在传递其 external_task_id 的任务
    • 并将其与您的协调器 DAG 的 execution_date 进行比较

    这应该澄清某些位

    【讨论】:

    • 感谢您的详细解释,感谢您的回复
    • @Kar 如果您能够解决问题,请考虑在此处添加解决方案作为答案,以供其他人参考。
    猜你喜欢
    • 1970-01-01
    • 2019-06-05
    • 2023-03-06
    • 1970-01-01
    • 2020-11-23
    • 1970-01-01
    • 2021-08-31
    • 2022-10-14
    • 1970-01-01
    相关资源
    最近更新 更多