כתיבת תרשימי DAG ב-Airflow

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

במדריך הזה נסביר איך לכתוב גרף אציקלי מכוון (DAG) של Apache Airflow שפועל בסביבת Managed Service for Apache Airflow.

מכיוון ש-Apache Airflow לא מספק בידוד חזק של DAG ומשימות, מומלץ להשתמש בסביבות ייצור ובדיקה נפרדות כדי למנוע הפרעות ב-DAG. מידע נוסף זמין במאמר בנושא בדיקת DAG.

מבנה של DAG ב-Airflow

גרף מכוון מחזורי (DAG) של Airflow מוגדר בקובץ Python ומורכב מהרכיבים הבאים:

  • הגדרת DAG
  • אופרטורים של זרימת אוויר
  • קשרי אופרטורים

בקטעי הקוד הבאים מוצגות דוגמאות לכל רכיב ללא הקשר.

הגדרת DAG

בדוגמה הבאה מוצגת הגדרה של Airflow DAG:

Airflow 2

import datetime

from airflow import models

default_dag_args = {
    # The start_date describes when a DAG is valid / can be run. Set this to a
    # fixed point in time rather than dynamically, since it is evaluated every
    # time a DAG is parsed. See:
    # https://airflow.apache.org/faq.html#what-s-the-deal-with-start-date
    "start_date": datetime.datetime(2018, 1, 1),
}

# Define a DAG (directed acyclic graph) of tasks.
# Any task you create within the context manager is automatically added to the
# DAG object.
with models.DAG(
    "composer_sample_simple_greeting",
    schedule_interval=datetime.timedelta(days=1),
    default_args=default_dag_args,
) as dag:

Airflow 1

import datetime

from airflow import models

default_dag_args = {
    # The start_date describes when a DAG is valid / can be run. Set this to a
    # fixed point in time rather than dynamically, since it is evaluated every
    # time a DAG is parsed. See:
    # https://airflow.apache.org/faq.html#what-s-the-deal-with-start-date
    "start_date": datetime.datetime(2018, 1, 1),
}

# Define a DAG (directed acyclic graph) of tasks.
# Any task you create within the context manager is automatically added to the
# DAG object.
with models.DAG(
    "composer_sample_simple_greeting",
    schedule_interval=datetime.timedelta(days=1),
    default_args=default_dag_args,
) as dag:

אופרטורים ומשימות

Airflow Operators מתארים את העבודה שצריך לבצע. משימה היא מופע ספציפי של אופרטור.

Airflow 2

from airflow.operators.bash import BashOperator
from airflow.operators.python import PythonOperator

    def greeting():
        import logging

        logging.info("Hello World!")

    # An instance of an operator is called a task. In this case, the
    # hello_python task calls the "greeting" Python function.
    hello_python = PythonOperator(task_id="hello", python_callable=greeting)

    # Likewise, the goodbye_bash task calls a Bash script.
    goodbye_bash = BashOperator(task_id="bye", bash_command="echo Goodbye.")

Airflow 1

from airflow.operators import bash_operator
from airflow.operators import python_operator

    def greeting():
        import logging

        logging.info("Hello World!")

    # An instance of an operator is called a task. In this case, the
    # hello_python task calls the "greeting" Python function.
    hello_python = python_operator.PythonOperator(
        task_id="hello", python_callable=greeting
    )

    # Likewise, the goodbye_bash task calls a Bash script.
    goodbye_bash = bash_operator.BashOperator(
        task_id="bye", bash_command="echo Goodbye."
    )

קשרים בין משימות

קשרי משימות מתארים את הסדר שבו צריך להשלים את העבודה.

Airflow 2

# Define the order in which the tasks complete by using the >> and <<
# operators. In this example, hello_python executes before goodbye_bash.
hello_python >> goodbye_bash

Airflow 1

# Define the order in which the tasks complete by using the >> and <<
# operators. In this example, hello_python executes before goodbye_bash.
hello_python >> goodbye_bash

דוגמה מלאה לתהליך עבודה של DAG ב-Python

תהליך העבודה הבא הוא תבנית DAG מלאה שפועלת ומורכבת משתי משימות: משימה hello_python ומשימה goodbye_bash:

Airflow 2


import datetime

from airflow import models

from airflow.operators.bash import BashOperator
from airflow.operators.python import PythonOperator



default_dag_args = {
    # The start_date describes when a DAG is valid / can be run. Set this to a
    # fixed point in time rather than dynamically, since it is evaluated every
    # time a DAG is parsed. See:
    # https://airflow.apache.org/faq.html#what-s-the-deal-with-start-date
    "start_date": datetime.datetime(2018, 1, 1),
}

# Define a DAG (directed acyclic graph) of tasks.
# Any task you create within the context manager is automatically added to the
# DAG object.
with models.DAG(
    "composer_sample_simple_greeting",
    schedule_interval=datetime.timedelta(days=1),
    default_args=default_dag_args,
) as dag:
    def greeting():
        import logging

        logging.info("Hello World!")

    # An instance of an operator is called a task. In this case, the
    # hello_python task calls the "greeting" Python function.
    hello_python = PythonOperator(task_id="hello", python_callable=greeting)

    # Likewise, the goodbye_bash task calls a Bash script.
    goodbye_bash = BashOperator(task_id="bye", bash_command="echo Goodbye.")

    # Define the order in which the tasks complete by using the >> and <<
    # operators. In this example, hello_python executes before goodbye_bash.
    hello_python >> goodbye_bash

Airflow 1


import datetime

from airflow import models

from airflow.operators import bash_operator
from airflow.operators import python_operator



default_dag_args = {
    # The start_date describes when a DAG is valid / can be run. Set this to a
    # fixed point in time rather than dynamically, since it is evaluated every
    # time a DAG is parsed. See:
    # https://airflow.apache.org/faq.html#what-s-the-deal-with-start-date
    "start_date": datetime.datetime(2018, 1, 1),
}

# Define a DAG (directed acyclic graph) of tasks.
# Any task you create within the context manager is automatically added to the
# DAG object.
with models.DAG(
    "composer_sample_simple_greeting",
    schedule_interval=datetime.timedelta(days=1),
    default_args=default_dag_args,
) as dag:
    def greeting():
        import logging

        logging.info("Hello World!")

    # An instance of an operator is called a task. In this case, the
    # hello_python task calls the "greeting" Python function.
    hello_python = python_operator.PythonOperator(
        task_id="hello", python_callable=greeting
    )

    # Likewise, the goodbye_bash task calls a Bash script.
    goodbye_bash = bash_operator.BashOperator(
        task_id="bye", bash_command="echo Goodbye."
    )

    # Define the order in which the tasks complete by using the >> and <<
    # operators. In this example, hello_python executes before goodbye_bash.
    hello_python >> goodbye_bash

מידע נוסף על הגדרת DAG ב-Airflow זמין במדריך ל-Airflow ובמאמר מושגים ב-Airflow.

אופרטורים של זרימת אוויר

בדוגמאות הבאות מוצגים כמה אופרטורים פופולריים של Airflow. לעיון במקור מידע מהימן על אופרטורים של Airflow, אפשר לעיין בחומר העזר בנושא אופרטורים ו-Hooks ובאינדקס הספקים.

BashOperator

משתמשים ב-BashOperator כדי להריץ תוכניות של שורת פקודה.

Airflow 2

from airflow.operators import bash

    # Create BigQuery output dataset.
    make_bq_dataset = bash.BashOperator(
        task_id="make_bq_dataset",
        # Executing 'bq' command requires Google Cloud SDK which comes
        # preinstalled in Cloud Composer.
        bash_command=f"bq ls {bq_dataset_name} || bq mk {bq_dataset_name}",
    )

Airflow 1

from airflow.operators import bash_operator

    # Create BigQuery output dataset.
    make_bq_dataset = bash_operator.BashOperator(
        task_id="make_bq_dataset",
        # Executing 'bq' command requires Google Cloud SDK which comes
        # preinstalled in Cloud Composer.
        bash_command=f"bq ls {bq_dataset_name} || bq mk {bq_dataset_name}",
    )

‫Managed Airflow מריץ את הפקודות שצוינו בסקריפט Bash ב-worker של Airflow. ה-worker הוא קונטיינר Docker מבוסס-Debian והוא כולל כמה חבילות.

PythonOperator

משתמשים ב-PythonOperator כדי להריץ קוד Python שרירותי.

‫Managed Airflow מריץ את קוד Python בקונטיינר שכולל חבילות של גרסת האימג' של Managed Airflow שבה נעשה שימוש בסביבה שלכם.

כדי להתקין חבילות Python נוספות, ראו התקנת יחסי תלות ב-Python.

Google Cloud אופרטורים

כדי להריץ משימות שמשתמשות במוצרי Google Cloud , משתמשים באופרטורים שלGoogle Cloud Airflow. לדוגמה, אופרטורים של BigQuery שולחים שאילתות ומעבדים נתונים ב-BigQuery.

יש עוד הרבה אופרטורים של Airflow ל Google Cloud ולשירותים ספציפיים שמוצעים על ידי Google Cloud. Google Cloud רשימה מלאה של אופרטורים

Airflow 2

from airflow.providers.google.cloud.operators import bigquery
from airflow.providers.google.cloud.transfers import bigquery_to_gcs

    bq_recent_questions_query = bigquery.BigQueryInsertJobOperator(
        task_id="bq_recent_questions_query",
        configuration={
            "query": {
                "query": RECENT_QUESTIONS_QUERY,
                "useLegacySql": False,
                "destinationTable": {
                    "projectId": project_id,
                    "datasetId": bq_dataset_name,
                    "tableId": bq_recent_questions_table_id,
                },
            }
        },
        location=location,
    )

Airflow 1

from airflow.contrib.operators import bigquery_operator

    # Query recent StackOverflow questions.
    bq_recent_questions_query = bigquery_operator.BigQueryOperator(
        task_id="bq_recent_questions_query",
        sql="""
        SELECT owner_display_name, title, view_count
        FROM `bigquery-public-data.stackoverflow.posts_questions`
        WHERE creation_date < CAST('{max_date}' AS TIMESTAMP)
            AND creation_date >= CAST('{min_date}' AS TIMESTAMP)
        ORDER BY view_count DESC
        LIMIT 100
        """.format(
            max_date=max_query_date, min_date=min_query_date
        ),
        use_legacy_sql=False,
        destination_dataset_table=bq_recent_questions_table_id,
    )

EmailOperator

משתמשים ב-EmailOperator כדי לשלוח אימייל מ-DAG. כדי לשלוח אימייל מסביבת Managed Airflow, צריך להגדיר את הסביבה לשימוש ב-SendGrid.

Airflow 2

from airflow.operators import email

    # Send email confirmation (you will need to set up the email operator
    # See https://cloud.google.com/composer/docs/how-to/managing/creating#notification
    # for more info on configuring the email operator in Cloud Composer)
    email_summary = email.EmailOperator(
        task_id="email_summary",
        to="{{var.value.email}}",
        subject="Sample BigQuery notify data ready",
        html_content="""
        Analyzed Stack Overflow posts data from {min_date} 12AM to {max_date}
        12AM. The most popular question was '{question_title}' with
        {view_count} views. Top 100 questions asked are now available at:
        {export_location}.
        """.format(
            min_date=min_query_date,
            max_date=max_query_date,
            question_title=(
                "{{ ti.xcom_pull(task_ids='bq_read_most_popular', "
                "key='return_value')[0][0] }}"
            ),
            view_count=(
                "{{ ti.xcom_pull(task_ids='bq_read_most_popular', "
                "key='return_value')[0][1] }}"
            ),
            export_location=output_file,
        ),
    )

Airflow 1

from airflow.operators import email_operator

    # Send email confirmation
    email_summary = email_operator.EmailOperator(
        task_id="email_summary",
        to="{{var.value.email}}",
        subject="Sample BigQuery notify data ready",
        html_content="""
        Analyzed Stack Overflow posts data from {min_date} 12AM to {max_date}
        12AM. The most popular question was '{question_title}' with
        {view_count} views. Top 100 questions asked are now available at:
        {export_location}.
        """.format(
            min_date=min_query_date,
            max_date=max_query_date,
            question_title=(
                "{{ ti.xcom_pull(task_ids='bq_read_most_popular', "
                "key='return_value')[0][0] }}"
            ),
            view_count=(
                "{{ ti.xcom_pull(task_ids='bq_read_most_popular', "
                "key='return_value')[0][1] }}"
            ),
            export_location=output_file,
        ),
    )

התראות על כשלים באופרטור

מגדירים את email_on_failure ל-True כדי לשלוח התראה באימייל כשמפעיל ב-DAG נכשל. כדי לשלוח התראות באימייל מסביבת Managed Airflow, צריך להגדיר את הסביבה לשימוש ב-SendGrid.

Airflow 2

from airflow import models

default_dag_args = {
    "start_date": yesterday,
    # Email whenever an Operator in the DAG fails.
    "email": "{{var.value.email}}",
    "email_on_failure": True,
    "email_on_retry": False,
    "retries": 1,
    "retry_delay": datetime.timedelta(minutes=5),
    "project_id": project_id,
}

with models.DAG(
    "composer_sample_bq_notify",
    schedule_interval=datetime.timedelta(weeks=4),
    default_args=default_dag_args,
) as dag:

Airflow 1

from airflow import models

default_dag_args = {
    "start_date": yesterday,
    # Email whenever an Operator in the DAG fails.
    "email": "{{var.value.email}}",
    "email_on_failure": True,
    "email_on_retry": False,
    "retries": 1,
    "retry_delay": datetime.timedelta(minutes=5),
    "project_id": "{{var.value.gcp_project}}",
}

with models.DAG(
    "composer_sample_bq_notify",
    schedule_interval=datetime.timedelta(weeks=4),
    default_args=default_dag_args,
) as dag:

הנחיות לתהליכי עבודה מסוג DAG

  • ממקמים את כל ספריות Python בהתאמה אישית בארכיון ZIP של DAG בספרייה מקוננת. אין למקם ספריות ברמה העליונה של ספריית ה-DAG.

    כש-Airflow סורק את התיקייה dags/, הוא בודק רק אם יש DAG במודולים של Python שנמצאים ברמה העליונה של תיקיית ה-DAG וברמה העליונה של ארכיון ZIP שנמצא גם הוא ברמה העליונה של התיקייה dags/. אם Airflow נתקל במודול Python בארכיון ZIP שלא מכיל את מחרוזות המשנה airflow ו-DAG, הוא מפסיק לעבד את ארכיון ה-ZIP. ‫Airflow מחזיר רק את ה-DAG שנמצאו עד לנקודה הזו.

  • משתמשים ב-Airflow 2 במקום ב-Airflow 1.

    קהילת Airflow לא מפרסמת יותר גרסאות משניות חדשות או גרסאות תיקון ל-Airflow 1.

  • כדי להבטיח עמידות בפני תקלות, אל תגדירו כמה אובייקטים של DAG באותו מודול Python.

  • אל תשתמשו ב-SubDAGs. במקום זאת, מקבצים משימות בתוך DAGs.

  • ממקמים את הקבצים הנדרשים בזמן הניתוח של DAG בתיקייה dags/, ולא בתיקייה data/.

  • הטמעה של בדיקות יחידה ל-DAG

  • מומלץ לבדוק DAG שפותח או שונה לפי ההוראות לבדיקת DAG.

  • מוודאים שפיתוח של DAG לא גורם לעלייה גדולה מדי בזמני הניתוח של DAG.

  • יכולות להיות כמה סיבות לכשל במשימות של Airflow. כדי למנוע כשלים בהפעלות של DAG שלם, מומלץ להפעיל ניסיונות חוזרים של משימות. הגדרה של מספר הניסיונות החוזרים המקסימלי לערך 0 פירושה שלא יתבצעו ניסיונות חוזרים.

    מומלץ לשנות את האפשרות default_task_retries לערך אחר של ניסיונות חוזרים לביצוע המשימה, ולא ל-0. בנוסף, אפשר להגדיר את הפרמטר retries ברמת המשימה.

  • אם רוצים להשתמש ב-GPU במשימות Airflow, צריך ליצור אשכול GKE נפרד שמבוסס על צמתים באמצעות מכונות עם GPU. משתמשים ב-GKEStartPodOperator כדי להריץ את המשימות.

  • מומלץ להימנע מהרצת משימות שצורכות הרבה משאבים של CPU וזיכרון במאגר הצמתים של האשכול, שבו פועלים רכיבים אחרים של Airflow (מתזמנים, עובדים, שרתי אינטרנט). במקום זאת, אפשר להשתמש ב-KubernetesPodOperator או ב-GKEStartPodOperator.

  • כשפורסים DAG בסביבה, מעלים רק את הקבצים שנדרשים לפרשנות ולביצוע של DAG בתיקייה /dags.

  • הגבלת מספר קובצי ה-DAG בתיקייה /dags.

    מערכת Airflow מנתחת באופן רציף קובצי DAG בתיקייה /dags. הניתוח הוא תהליך שחוזר על עצמו בתיקיית ה-DAG, ומספר הקבצים שצריך לטעון (עם התלות שלהם) משפיע על הביצועים של ניתוח ה-DAG ותזמון המשימות. יעיל הרבה יותר להשתמש ב-100 קבצים עם 100 DAG בכל אחד מאשר ב-10,000 קבצים עם DAG אחד בכל אחד, ולכן מומלץ לבצע אופטימיזציה כזו. האופטימיזציה הזו היא איזון בין זמן הניתוח לבין היעילות של יצירה וניהול של DAG.

    לדוגמה, אם רוצים לפרוס 10,000 קובצי DAG, אפשר ליצור 100 קובצי ZIP, שכל אחד מהם מכיל 100 קובצי DAG.

    בנוסף להצעות שלמעלה, אם יש לכם יותר מ-10,000 קובצי DAG, יכול להיות שכדאי ליצור DAG באופן פרוגרמטי. לדוגמה, אתם יכולים להטמיע קובץ Python DAG יחיד שיוצר מספר מסוים של אובייקטים של DAG (לדוגמה, 20 או 100 אובייקטים של DAG).

  • לא מומלץ להשתמש באופרטורים של Airflow שהוצאו משימוש. במקום זאת, אפשר להשתמש בחלופות העדכניות שלהם.

  • לא מומלץ להשתמש ב-CloudLoggingHandler לטיפול ברישום ביומן. ה-handler הזה של הרישום גורם לעומס יתר במסד הנתונים של הסביבה ולכשלים הבאים ב-DAG עם השגיאות 'השרת סגר את החיבור באופן לא צפוי'.

שאלות נפוצות בנושא כתיבת DAG

איך אפשר לצמצם את כפילויות הקוד אם רוצים להריץ את אותן משימות או משימות דומות בכמה DAG?

מומלץ להגדיר ספריות ועטיפות כדי לצמצם את כפילות הקוד.

איך אפשר לעשות שימוש חוזר בקוד בין קובצי DAG?

ממקמים את פונקציות השירות בספריית Python מקומית ומייבאים את הפונקציות. אפשר להפנות לפונקציות בכל DAG שנמצא בתיקייה dags/ בדלי של הסביבה.

איך אפשר לצמצם את הסיכון להגדרות שונות?

לדוגמה, יש לכם שני צוותים שרוצים לצבור נתונים גולמיים למדדי הכנסות. הצוותים כותבים שתי משימות שונות במקצת שמשיגות את אותו הדבר. הגדרת ספריות כדי לעבוד עם נתוני ההכנסות, כך שהמטמיעים של DAG צריכים להבהיר את ההגדרה של ההכנסות שמצטברות.

איך מגדירים תלות בין DAG?

ההגדרה תלויה באופן שבו רוצים להגדיר את התלות.

אם יש לכם שני DAG‏ (DAG A ו-DAG B) ואתם רוצים ש-DAG B יופעל אחרי DAG A, אתם יכולים להוסיף TriggerDagRunOperator בסוף DAG A.

אם DAG B תלוי רק בארטיפקט שנוצר על ידי DAG A, כמו הודעת Pub/Sub, יכול להיות שחיישן יתאים יותר.

אם DAG B משולב בצורה הדוקה עם DAG A, יכול להיות שאפשר למזג את שני ה-DAG ל-DAG אחד.

איך מעבירים מזהי הפעלה ייחודיים ל-DAG ולמשימות שלו?

לדוגמה, אתם רוצים להעביר שמות של אשכולות Dataproc ונתיבי קבצים.

כדי ליצור מזהה ייחודי אקראי, מחזירים str(uuid.uuid4()) ב-PythonOperator. המזהה יופיע ב-XComs, כך שתוכלו להפנות למזהה באופרטורים אחרים באמצעות שדות של תבניות.

לפני שמפיקים uuid, כדאי לשקול אם מזהה ספציפי של DagRun יהיה שימושי יותר. אפשר גם להפנות למזהים האלה בהחלפות של Jinja באמצעות פקודות מאקרו.

איך מפרידים בין משימות ב-DAG?

כל משימה צריכה להיות יחידת עבודה אידמפוטנטית. לכן, מומלץ להימנע מהוספת תהליך עבודה מרובה שלבים בתוך משימה אחת, כמו תוכנית מורכבת שפועלת ב-PythonOperator.

האם כדאי להגדיר כמה משימות ב-DAG אחד כדי לצבור נתונים מכמה מקורות?

לדוגמה, יש לכם כמה טבלאות עם נתונים גולמיים ואתם רוצים ליצור נתונים מצטברים יומיים לכל טבלה. המשימות לא תלויות זו בזו. כדאי ליצור משימה אחת ו-DAG אחד לכל טבלה, או ליצור DAG כללי אחד?

אם אתם רוצים שכל משימה תשתף את אותם מאפיינים ברמת ה-DAG, כמו schedule, כדאי להגדיר כמה משימות ב-DAG אחד. אחרת, כדי לצמצם את כפילות הקוד, אפשר ליצור כמה DAG ממודול Python יחיד על ידי הצבתם ב-globals() של המודול.

איך מגבילים את מספר המשימות שפועלות בו-זמנית ב-DAG?

לדוגמה, אם אתם רוצים להימנע מחריגה ממכסות או ממגבלות שימוש ב-API, או מריצה של יותר מדי תהליכים בו-זמנית.

אתם יכולים להגדיר מאגרי Airflow בממשק המשתמש האינטרנטי של Airflow ולשייך משימות למאגרים קיימים ב-DAG.

שאלות נפוצות על שימוש באופרטורים

כדאי להשתמש ב-DockerOperator?

אנחנו לא ממליצים להשתמש ב-DockerOperator, אלא אם משתמשים בו להפעלת קונטיינרים בהתקנת Docker מרחוק (לא בתוך אשכול של סביבה). בסביבת Managed Airflow, לאופרטור אין גישה לדמונים של Docker.

במקום זאת, צריך להשתמש ב-KubernetesPodOperator או ב-GKEStartPodOperator. האופרטורים האלה מפעילים pods של Kubernetes באשכולות Kubernetes או GKE, בהתאמה. שימו לב: אנחנו לא ממליצים להפעיל pods באשכול של סביבה, כי זה עלול להוביל לתחרות על משאבים.

כדאי להשתמש ב-SubDagOperator?

לא מומלץ להשתמש ב-SubDagOperator.

כדאי להשתמש בחלופות כמו שמפורט במאמר קיבוץ משימות.

האם כדאי להריץ קוד Python רק ב-PythonOperators כדי להפריד באופן מלא בין אופרטורים של Python?

בהתאם ליעד, יש לכם כמה אפשרויות.

אם הדבר היחיד שחשוב לכם הוא לשמור על יחסי תלות נפרדים של Python, אתם יכולים להשתמש ב-PythonVirtualenvOperator.

כדאי להשתמש ב-KubernetesPodOperator. האופרטור הזה מאפשר להגדיר קבוצות Pod של Kubernetes ולהריץ אותן באשכולות אחרים.

איך מוסיפים חבילות בינאריות או חבילות שאינן PyPI בהתאמה אישית?

אתם יכולים להתקין חבילות שמתארחות במאגרי חבילות פרטיים.

איך מעבירים ארגומנטים באופן אחיד ל-DAG ולמשימות שלו?

אתם יכולים להשתמש בתמיכה המובנית של Airflow בתבניות Jinja כדי להעביר ארגומנטים שאפשר להשתמש בהם בשדות של תבניות.

מתי מתבצעת החלפת התבנית?

החלפת התבנית מתבצעת בתהליכי העבודה של Airflow ממש לפני שקוראים לפונקציה pre_execute של אופרטור. בפועל, המשמעות היא שהתבניות לא מוחלפות עד לרגע שלפני הפעלת המשימה.

איך אפשר לדעת אילו ארגומנטים של אופרטור תומכים בהחלפת תבנית?

ארגומנטים של אופרטורים שתומכים בהחלפת תבניות Jinja2 מסומנים במפורש ככאלה.

מחפשים את השדה template_fields בהגדרת האופרטור, שמכיל רשימה של שמות ארגומנטים שעוברים החלפה בתבנית.

לדוגמה, אפשר לעיין ב-BashOperator, שתומך ביצירת תבניות לארגומנטים bash_command ו-env.

אופרטורים של Airflow שהוצאו משימוש והוסרו

האופרטורים של Airflow שמפורטים בטבלה הבאה הוצאו משימוש:

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

  • בגרסת התחזוקה האחרונה של Managed Airflow (Legacy Gen 1) (1.20.12) נעשה שימוש בגרסה של חבילת apache-airflow-providers-google שמבוססת על גרסה ציבורית 8.9.0. בטבלה הבאה אפשר לראות אם אופרטור מסוים נתמך או אם יש לו תחליף, בהתאם לגרסה הרשומה של החבילה.

  • חלק מהאופרטורים להחלפה לא נתמכים באף גרסה של Managed Airflow (דור קודם 1). כדי להשתמש בהם, כדאי לשדרג ל-Managed Airflow (דור 3) או ל-Managed Airflow (דור 2).

אופרטור שהוצא משימוש אופרטור החלפה

GCSObjectsWtihPrefixExistenceSensor

GCSObjectsWithPrefixExistenceSensor

CreateAutoMLTextTrainingJobOperator

SupervisedFineTuningTrainOperator

GKEDeploymentHook

GKEKubernetesHook

GKECustomResourceHook

GKEKubernetesHook

GKEPodHook

GKEKubernetesHook

GKEJobHook

GKEKubernetesHook

GKEPodAsyncHook

GKEKubernetesAsyncHook

SecretsManagerHook

GoogleCloudSecretManagerHook

BigQueryExecuteQueryOperator

BigQueryInsertJobOperator

BigQueryPatchDatasetOperator

BigQueryUpdateDatasetOperator

DataflowCreateJavaJobOperator

beam.BeamRunJavaPipelineOperator

DataflowCreatePythonJobOperator

beam.BeamRunPythonPipelineOperator

DataprocSubmitPigJobOperator

DataprocSubmitJobOperator

DataprocSubmitHiveJobOperator

DataprocSubmitJobOperator

DataprocSubmitSparkSqlJobOperator

DataprocSubmitJobOperator

DataprocSubmitSparkJobOperator

DataprocSubmitJobOperator

DataprocSubmitHadoopJobOperator

DataprocSubmitJobOperator

DataprocSubmitPySparkJobOperator

DataprocSubmitJobOperator

BigQueryTableExistenceAsyncSensor

BigQueryTableExistenceSensor

BigQueryTableExistencePartitionAsyncSensor

BigQueryTablePartitionExistenceSensor

CloudComposerEnvironmentSensor

CloudComposerCreateEnvironmentOperator, ‏ CloudComposerDeleteEnvironmentOperator, ‏ CloudComposerUpdateEnvironmentOperator

GCSObjectExistenceAsyncSensor

GCSObjectExistenceSensor

GoogleAnalyticsHook

GoogleAnalyticsAdminHook

GoogleAnalyticsListAccountsOperator

GoogleAnalyticsAdminListAccountsOperator

GoogleAnalyticsGetAdsLinkOperator

GoogleAnalyticsAdminGetGoogleAdsLinkOperator

GoogleAnalyticsRetrieveAdsLinksListOperator

GoogleAnalyticsAdminListGoogleAdsLinksOperator

GoogleAnalyticsDataImportUploadOperator

GoogleAnalyticsAdminCreateDataStreamOperator

GoogleAnalyticsDeletePreviousDataUploadsOperator

GoogleAnalyticsAdminDeleteDataStreamOperator

DataPipelineHook

DataflowHook

CreateDataPipelineOperator

DataflowCreatePipelineOperator

RunDataPipelineOperator

DataflowRunPipelineOperator

AutoMLBatchPredictOperator

vertex_ai.batch_prediction_job

PromptLanguageModelOperator

TextGenerationModelPredictOperator

GenerateTextEmbeddingsOperator

TextEmbeddingModelGetEmbeddingsOperator

PromptMultimodalModelOperator

GenerativeModelGenerateContentOperator

PromptMultimodalModelWithMediaOperator

GenerativeModelGenerateContentOperator

DataflowStartSqlJobOperator

DataflowStartYamlJobOperator

DataprocScaleClusterOperator

DataprocUpdateClusterOperator

MLEngineStartBatchPredictionJobOperator

CreateBatchPredictionJobOperator

MLEngineManageModelOperator

MLEngineCreateModelOperator, MLEngineGetModelOperator

MLEngineGetModelOperator

GetModelOperator

MLEngineDeleteModelOperator

DeleteModelOperator

MLEngineManageVersionOperator

MLEngineCreateVersion, MLEngineSetDefaultVersion, MLEngineListVersions, MLEngineDeleteVersion

MLEngineCreateVersionOperator

פרמטר parent_model לאופרטורים של VertexAI

MLEngineSetDefaultVersionOperator

SetDefaultVersionOnModelOperator

MLEngineListVersionsOperator

ListModelVersionsOperator

MLEngineDeleteVersionOperator

DeleteModelVersionOperator

MLEngineStartTrainingJobOperator

CreateCustomPythonPackageTrainingJobOperator

MLEngineTrainingCancelJobOperator

CancelCustomTrainingJobOperator

AutoMLDatasetLink

TranslationLegacyDatasetLink

AutoMLDatasetListLink

TranslationDatasetListLink

AutoMLModelLink

TranslationLegacyModelLink

AutoMLModelTrainLink

TranslationLegacyModelTrainLink

AutoMLModelPredictLink

TranslationLegacyModelPredictLink

TextGenerationModelPredictOperator

GenerativeModelGenerateContentOperator

LifeSciencesHook

ה-hook של Google Cloud Batch Operators

LifeSciencesLink

אין תחליף (ה-API הוצא משימוש)

LifeSciencesRunPipelineOperator

Google Cloud Batch Operators

BigQueryCreateEmptyTableOperator

BigQueryCreateTableOperator

BigQueryCreateExternalTableOperator

BigQueryCreateTableOperator

GoogleDisplayVideo360CreateQueryOperator

display_video.GoogleDisplayVideo360CreateSDFDownloadTaskOperator

GoogleDisplayVideo360RunQueryOperator

display_video.GoogleDisplayVideo360CreateSDFDownloadTaskOperator

GoogleDisplayVideo360DeleteReportOperator

משימת ייצוא של קובץ SDF בגרסה 4 של ה-API

GoogleDisplayVideo360DownloadReportV2Operator

display_video.GoogleDisplayVideo360SDFtoGCSOperator

GoogleDisplayVideo360DownloadLineItemsOperator

display_video.GoogleDisplayVideo360CreateSDFDownloadTaskOperator

GoogleDisplayVideo360UploadLineItemsOperator

display_video.GoogleDisplayVideo360SDFtoGCSOperator

GoogleDisplayVideo360RunQuerySensor

משימת ייצוא של קובץ SDF בגרסה 4 של ה-API

AutoMLPredictOperator

‫TextGenerationModelPredictOperator, ‏ TranslateTextOperator

CloudAutoMLHook

vertex_ai.auto_ml.AutoMLHook, translate.TranslateHook

AutoMLCreateDatasetOperator

vertex_ai.dataset.CreateDatasetOperator, translate.TranslateCreateDatasetOperator

AutoMLDeleteDatasetOperator

vertex_ai.dataset.ListDatasetsOperator, translate.TranslateDatasetsListOperator

AutoMLDeleteModelOperator

vertex_ai.model_service.DeleteModelOperator, translate.TranslateDeleteModelOperator

AutoMLDeployModelOperator

vertex_ai.endpoint_service.DeployModelOperator

AutoMLGetModelOperator

vertex_ai.model_service.GetModelOperator

AutoMLImportDataOperator

vertex_ai.dataset.ImportDataOperator, translate.TranslateImportDataOperator

AutoMLListDatasetOperator

vertex_ai.dataset.ListDatasetsOperator, translate.TranslateDatasetsListOperator

AutoMLTablesListColumnSpecsOperator

אין תחליף (הגרסה הקודמת של AutoML Tables הושבתה)

AutoMLTablesListTableSpecsOperator

אין תחליף (הגרסה הקודמת של AutoML Tables הושבתה)

AutoMLTablesUpdateDatasetOperator

vertex_ai.dataset.UpdateDatasetOperator

AutoMLTrainModelOperator

vertex_ai.auto_ml.CreateAutoMLTabularTrainingJobOperator, CreateAutoMLVideoTrainingJobOperator, CreateAutoMLImageTrainingJobOperator, SupervisedFineTuningTrainOperator, TranslateCreateModelOperator

MLEngineCreateModelOperator

אופרטור VertexAI תואם

DataCatalogEntryGroupLink

DataplexCatalogEntryGroupLink

DataCatalogEntryLink

DataplexCatalogEntryLink

DataCatalogTagTemplateLink

DataplexCatalogAspectTypeLink

CloudDataCatalogCreateEntryOperator

DataplexCatalogCreateEntryOperator

CloudDataCatalogCreateEntryGroupOperator

DataplexCatalogCreateEntryGroupOperator

CloudDataCatalogCreateTagOperator

DataplexCatalogCreateEntryOperator, ‏ DataplexCatalogUpdateEntryOperator

CloudDataCatalogCreateTagTemplateOperator

DataplexCatalogCreateAspectTypeOperator

CloudDataCatalogCreateTagTemplateFieldOperator

‫DataplexCatalogUpdateAspectTypeOperator, ‏ DataplexCatalogCreateAspectTypeOperator

CloudDataCatalogDeleteEntryGroupOperator

DataplexCatalogDeleteEntryGroupOperator

CloudDataCatalogDeleteTagOperator

DataplexCatalogUpdateEntryOperator

CloudDataCatalogDeleteTagTemplateOperator

DataplexCatalogDeleteAspectTypeOperator

CloudDataCatalogDeleteTagTemplateFieldOperator

DataplexCatalogUpdateAspectTypeOperator

CloudDataCatalogGetEntryOperator

DataplexCatalogGetEntryOperator

CloudDataCatalogGetEntryGroupOperator

DataplexCatalogGetEntryGroupOperator

CloudDataCatalogGetTagTemplateOperator

DataplexCatalogGetAspectTypeOperator

CloudDataCatalogListTagsOperator

DataplexCatalogGetEntryOperator

CloudDataCatalogLookupEntryOperator

DataplexCatalogLookupEntryOperator

CloudDataCatalogRenameTagTemplateFieldOperator

DataplexCatalogUpdateAspectTypeOperator

CloudDataCatalogSearchCatalogOperator

DataplexCatalogSearchEntriesOperator

CloudDataCatalogUpdateEntryOperator

DataplexCatalogUpdateEntryOperator

CloudDataCatalogUpdateTagOperator

DataplexCatalogUpdateEntryOperator

CloudDataCatalogUpdateTagTemplateOperator

DataplexCatalogUpdateAspectTypeOperator

CloudDataCatalogUpdateTagTemplateFieldOperator

DataplexCatalogUpdateAspectTypeOperator

CloudDataCatalogHook

DataplexHook

TextEmbeddingModelGetEmbeddingsOperator

GenAIGenerateEmbeddingsOperator

GenerativeModelGenerateContentOperator

GenAIGenerateContentOperator

SupervisedFineTuningTrainOperator

GenAISupervisedFineTuningTrainOperator

CountTokensOperator

GenAICountTokensOperator

CreateCachedContentOperator

GenAICreateCachedContentOperator

GenerateFromCachedContentOperator

GenAIGenerateContentOperator

DeleteExperimentRunOperator

DeleteExperimentRunOperator

TranslationLegacyModelLink

TranslationModelLink

TranslationDatasetListLink

TranslationDatasetsListLink

TranslationLegacyDatasetLink

TranslationDatasetsListLink

TranslationLegacyModelTrainLink

TranslationModelsListLink

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