【问题标题】:Issues with importing airflow.operators.sensors.external_task import ExternalTaskSensor module and triggering external dag导入 airflow.operators.sensors.external_task 导入 ExternalTask​​Sensor 模块和触发外部 dag 的问题
【发布时间】:2021-09-06 12:07:02
【问题描述】:

我正在尝试通过 master dag 触发多个外部 dag 数据流作业。

我打算使用 TriggerDagRunOperator 和 ExternalTask​​Sensor 。我有大约 10 个数据流作业 - 有些是按顺序执行的,有些是并行执行的。 例如:我想从 master dag 执行 Dag 数据流作业 A、B、C 等,在执行下一个任务之前,我想确保之前的 dag 运行已经完成。但我在导入 ExternalTask​​Sensor 模块时遇到问题。 他们是否有任何替代途径来实现这一目标?

注意:每个 DAG eg A/B/C 有 6-7 个任务。ExternalTask​​Sensor 可以在 DAG B 或 C 开始之前检查 dag A 的最后一个任务是否完成。

【问题讨论】:

  • 嗨@recyclinguy,建议在提问时提供一些示例代码。查看如何ask 提问。您能否提供在导入 ExternalTask​​Sensor 模块时收到的错误消息。

标签: airflow google-cloud-composer


【解决方案1】:

我使用下面的示例代码来运行使用 ExternalTask​​Sensor 的 dag,我能够成功导入 ExternalTask​​Sensor 模块。

import time
from datetime import datetime, timedelta
from pprint import pprint

from airflow import DAG
from airflow.operators.dagrun_operator import TriggerDagRunOperator
from airflow.operators.dummy_operator import DummyOperator
from airflow.operators.python_operator import PythonOperator
from airflow.sensors.external_task_sensor import ExternalTaskSensor
from airflow.utils.state import State

sensors_dag = DAG(
    "test_launch_sensors",
    schedule_interval=None,
    start_date=datetime(2020, 2, 14, 0, 0, 0),
    dagrun_timeout=timedelta(minutes=150),
    tags=["DEMO"],
)

dummy_dag = DAG(
    "test_dummy_dag",
    schedule_interval=None,
    start_date=datetime(2020, 2, 14, 0, 0, 0),
    dagrun_timeout=timedelta(minutes=150),
    tags=["DEMO"],
)


def print_context(ds, **context):
    pprint(context['conf'])


with dummy_dag:
    starts = DummyOperator(task_id="starts", dag=dummy_dag)
    empty = PythonOperator(
        task_id="empty",
        provide_context=True,
        python_callable=print_context,
        dag=dummy_dag,
    )
    ends = DummyOperator(task_id="ends", dag=dummy_dag)

    starts >> empty >> ends

with sensors_dag:
    trigger = TriggerDagRunOperator(
        task_id=f"trigger_{dummy_dag.dag_id}",
        trigger_dag_id=dummy_dag.dag_id,
        conf={"key": "value"},
        execution_date="{{ execution_date }}",
    )
    sensor = ExternalTaskSensor(
        task_id="wait_for_dag",
        external_dag_id=dummy_dag.dag_id,
        external_task_id="ends",
        poke_interval=5,
        timeout=120,
    )
    trigger >> sensor

在上述示例代码中,sensors_dag 使用 TriggerDagRunOperator() 触发 dummy_dag 中的任务。 sensor_dag 会一直等到 dummy_dag 中指定的 external_task 完成。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-06-05
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-06-29
    • 1970-01-01
    相关资源
    最近更新 更多