Managed Airflow (דור 3) | Managed Airflow (דור 2) | Managed Airflow (דור 1 מדור קודם)
בדף הזה נדגים איך מטמיעים DAG שמפעיל DAGs בסביבות ובפרויקטים אחרים של Managed Airflow באמצעות אופרטורים של Airflow ל-Managed Airflow.
אם רוצים להפעיל DAG בסביבה, אפשר לעיין במאמר בנושא תזמון והפעלה של DAG.
הגדרת הרשאות IAM
אם סביבת היעד נמצאת בפרויקט אחר, לחשבון השירות של הסביבה שלכם צריכות להיות הרשאות שמאפשרות אינטראקציה עם סביבות באותו פרויקט.
| פרויקט | משאב | חשבון משתמש | תפקיד |
|---|---|---|---|
| הפרויקט שבו נמצאת סביבת היעד | פרויקט | חשבון השירות של סביבת המקור |
התפקיד Composer Worker (composer.worker) |
| הפרויקט שבו נמצאת סביבת היעד | פרויקט | חשבון השירות של סביבת המקור |
תפקיד מותאם אישית עם ההרשאה composer.environments.executeAirflowCommand |
הפעלת DAG בסביבה אחרת
ה-DAG לדוגמה שמתואר בקטע הזה מבצע את הפעולות הבאות:
- הפעלת DAG בסביבת Managed Airflow אחרת.
- בודק אם הפעלת DAG הסתיימה.
אחרי שהרצת ה-DAG בסביבה אחרת מסתיימת, ה-DAG לדוגמה מסומן כהצלחה.
הרצת פקודות Airflow CLI באמצעות CloudComposerRunAirflowCLICommandOperator
אפשר להשתמש באופרטור CloudComposerRunAirflowCLICommandOperator כדי להריץ פקודות Airflow CLI בסביבת Managed Airflow אחרת. ה-DAG לדוגמה מריץ את הפקודה dags trigger, שמפעילה DAG.
אפשר להפעיל את האופרטור הזה במצב ניתן להשהיה על ידי הגדרת הפרמטר deferrable לערך True.
run_airflow_cli_cmd = CloudComposerRunAirflowCLICommandOperator(
task_id="run_airflow_cli_cmd",
project_id="target-project",
environment_id="target-composer-environment",
region="us-central1",
command="dags trigger -- target_dag",
# You can run this operator in the deferrable mode:
# deferrable=True
)
איך בודקים אם הפעלת DAG הסתיימה
אפשר להשתמש בחיישן CloudComposerDAGRunSensor כדי לבדוק אם הפעלת DAG הושלמה בסביבת Managed Airflow אחרת.
החיישן הזה יכול לפעול במצב דחייה. כדי להפעיל אותו, צריך להגדיר את הפרמטר deferrable לערך True.
dag_run_sensor = CloudComposerDAGRunSensor(
task_id="dag_run_sensor",
project_id="target-project",
environment_id="target-composer-environment",
region="us-central1",
composer_dag_id="target_dag",
allowed_states=["success"],
# You can run this sensor in the deferrable mode:
# deferrable=True
)
קוד מלא לדוגמה
זוהי דוגמת קוד מלאה של DAG שמשלב את שתי המשימות שתוארו קודם.
from datetime import datetime, timedelta
from airflow.models.dag import DAG
from airflow.providers.google.cloud.operators.cloud_composer import (
CloudComposerRunAirflowCLICommandOperator,
)
from airflow.providers.google.cloud.sensors.cloud_composer import CloudComposerDAGRunSensor
DAG_ID = "trigger_dag_in_another_composer_environment"
TARGET_PROJECT_ID = "example-target-project"
TARGET_REGION = "example-target-region"
TARGET_ENV_ID = "example-target-composer-environment"
TARGET_DAG = "example_target_dag_id"
COMMAND = f"dags trigger -- {TARGET_DAG}"
with DAG(
DAG_ID,
schedule="@once",
start_date=datetime(2024, 1, 1),
catchup=False,
tags=["example", "composer"],
) as dag:
run_airflow_cli_cmd = CloudComposerRunAirflowCLICommandOperator(
task_id="run_airflow_cli_cmd",
project_id=TARGET_PROJECT_ID,
environment_id=TARGET_ENV_ID,
region=TARGET_REGION,
command=COMMAND,
# You can run this operator in the deferrable mode:
# deferrable=True
)
dag_run_sensor = CloudComposerDAGRunSensor(
task_id="dag_run_sensor",
project_id=TARGET_PROJECT_ID,
environment_id=TARGET_ENV_ID,
region=TARGET_REGION,
composer_dag_id=TARGET_DAG,
allowed_states=["success"],
execution_range=timedelta(minutes=5),
# You can run this sensor in the deferrable mode:
# deferrable=True
)
run_airflow_cli_cmd >> dag_run_sensor