קיבוץ משימות בתוך DAG

Managed Airflow (דור 3) | Managed Airflow (דור 2) | Managed Airflow (דור 1 מדור קודם)

בדף הזה מוסבר איך לקבץ משימות בצינורות של Airflow באמצעות דפוסי העיצוב הבאים:

  • קיבוץ משימות בתרשים DAG.
  • הפעלת DAGs צאצאים מ-DAG הורה.
  • קיבוץ משימות באמצעות האופרטור TaskGroup.

קיבוץ משימות בתרשים DAG

כדי לקבץ משימות בשלבים מסוימים של צינור הנתונים, אפשר להשתמש בקשרים בין המשימות בקובץ ה-DAG.

דוגמה:

תרשים של משימות Airflow שמציג משימות הסתעפות
איור 1. אפשר לקבץ משימות ב-DAG של Airflow (לחיצה להגדלה)

בתהליך העבודה הזה, המשימות op-1 ו-op-2 פועלות יחד אחרי המשימה הראשונית start. כדי לעשות את זה, מקבצים את המשימות באמצעות ההצהרה start >> [task_1, task_2].

בדוגמה הבאה מוצגת הטמעה מלאה של ה-DAG הזה:

from airflow import DAG
from airflow.operators.bash import BashOperator
from airflow.operators.dummy import DummyOperator
from airflow.utils.dates import days_ago

DAG_NAME = "all_tasks_in_one_dag"

args = {"owner": "airflow", "start_date": days_ago(1), "schedule_interval": "@once"}

with DAG(dag_id=DAG_NAME, default_args=args) as dag:
    start = DummyOperator(task_id="start")

    task_1 = BashOperator(task_id="op-1", bash_command=":", dag=dag)

    task_2 = BashOperator(task_id="op-2", bash_command=":", dag=dag)

    some_other_task = DummyOperator(task_id="some-other-task")

    task_3 = BashOperator(task_id="op-3", bash_command=":", dag=dag)

    task_4 = BashOperator(task_id="op-4", bash_command=":", dag=dag)

    end = DummyOperator(task_id="end")

    start >> [task_1, task_2] >> some_other_task >> [task_3, task_4] >> end

הפעלת DAGs צאצאים מ-DAG אב

אפשר להפעיל DAG אחד מ-DAG אחר באמצעות האופרטור TriggerDagRunOperator.

דוגמה:

תרשים של משימות Airflow שמציג DAGs של צאצאים שהופעלו כחלק מתרשים DAG
איור 2. אפשר להפעיל DAG מתוך DAG באמצעות TriggerDagRunOperator (לוחצים להגדלה)

בתרשים זרימת העבודה הזה, הבלוקים dag_1 ו-dag_2 מייצגים סדרה של משימות שמקובצות יחד ב-DAG נפרד בסביבת Managed Airflow.

ההטמעה של תהליך העבודה הזה מחייבת שני קובצי DAG נפרדים. קובץ ה-DAG של הבקרה נראה כך:

from airflow import DAG
from airflow.operators.dummy import DummyOperator
from airflow.operators.trigger_dagrun import TriggerDagRunOperator
from airflow.utils.dates import days_ago


with DAG(
    dag_id="controller_dag_to_trigger_other_dags",
    default_args={"owner": "airflow"},
    start_date=days_ago(1),
    schedule_interval="@once",
) as dag:
    start = DummyOperator(task_id="start")

    trigger_1 = TriggerDagRunOperator(
        task_id="dag_1",
        trigger_dag_id="dag-to-trigger",  # Ensure this equals the dag_id of the DAG to trigger
        conf={"message": "Hello World"},
    )
    trigger_2 = TriggerDagRunOperator(
        task_id="dag_2",
        trigger_dag_id="dag-to-trigger",  # Ensure this equals the dag_id of the DAG to trigger
        conf={"message": "Hello World"},
    )

    some_other_task = DummyOperator(task_id="some-other-task")

    end = DummyOperator(task_id="end")

    start >> trigger_1 >> some_other_task >> trigger_2 >> end

היישום של ה-DAG המשני, שמופעל על ידי ה-DAG הראשי, נראה כך:

from airflow import DAG
from airflow.operators.dummy import DummyOperator
from airflow.utils.dates import days_ago

DAG_NAME = "dag-to-trigger"

args = {"owner": "airflow", "start_date": days_ago(1), "schedule_interval": "None"}

with DAG(dag_id=DAG_NAME, default_args=args) as dag:
    dag_task = DummyOperator(task_id="dag-task")

כדי ש-DAG יפעל, צריך להעלות את שני קובצי ה-DAG בסביבת Managed Airflow.

קיבוץ משימות באמצעות האופרטור TaskGroup

אפשר להשתמש באופרטור TaskGroup כדי לקבץ משימות ב-DAG. משימות שמוגדרות בתוך בלוק TaskGroup עדיין נכללות ב-DAG הראשי.

דוגמה:

תרשים של משימות Airflow שמציג שתי קבוצות משימות
איור 3. אפשר לקבץ משימות באופן חזותי בממשק המשתמש באמצעות האופרטור TaskGroup (לחצו להגדלה)

המשימות op-1 ו-op-2 מקובצות יחד בבלוק עם המזהה taskgroup_1. הטמעה של תהליך העבודה הזה נראית כמו הקוד הבא:

from airflow.models.dag import DAG
from airflow.operators.bash import BashOperator
from airflow.operators.dummy import DummyOperator
from airflow.utils.dates import days_ago
from airflow.utils.task_group import TaskGroup

with DAG(dag_id="taskgroup_example", start_date=days_ago(1)) as dag:
    start = DummyOperator(task_id="start")

    with TaskGroup("taskgroup_1", tooltip="task group #1") as section_1:
        task_1 = BashOperator(task_id="op-1", bash_command=":")
        task_2 = BashOperator(task_id="op-2", bash_command=":")

    with TaskGroup("taskgroup_2", tooltip="task group #2") as section_2:
        task_3 = BashOperator(task_id="op-3", bash_command=":")
        task_4 = BashOperator(task_id="op-4", bash_command=":")

    some_other_task = DummyOperator(task_id="some-other-task")

    end = DummyOperator(task_id="end")

    start >> section_1 >> some_other_task >> section_2 >> end

המאמרים הבאים