Wyzwalanie podrzędnego DAGa
Zauważasz, że niektóre z twoich przepływów pracy korzystają z podobnych komponentów i dochodzisz do wniosku, że wspólne zadania można wydzielić do osobnego DAGa. Dzięki temu będzie można uruchamiać te komponenty w razie potrzeby, bez konieczności utrzymywania wielu ich kopii. Postanawiasz wywołać podrzędny DAG jako zadanie w bieżącym przepływie pracy.
Komponenty dag, task i datetime zostały już zaimportowane.
To ćwiczenie jest częścią kursu
Wprowadzenie do Apache Airflow w Pythonie
Instrukcje do ćwiczenia
- Zaimportuj operator potrzebny do uruchomienia DAGa z poziomu bieżącego przepływu pracy.
- Skonfiguruj operator tak, aby wyzwalał DAGa o nazwie
child_pipeline. - Upewnij się, że nadrzędny DAG czeka na zakończenie wyzwolonego DAGa przed kontynuowaniem.
- Ustaw, jak często operator sprawdza, czy podrzędny DAG zakończył działanie.
Interaktywne ćwiczenie praktyczne
Spróbuj tego ćwiczenia, uzupełniając ten przykładowy kod.
# Import TriggerDagRunOperator
from airflow.providers.standard.operators.trigger_dagrun import ____
@dag(start_date=datetime(2026, 1, 1))
def parent_orchestrator_dag():
# Trigger child_pipeline and wait for it to complete
trigger_child = TriggerDagRunOperator(
task_id="trigger_child_pipeline",
trigger_dag_id="____",
____=True,
____=30,
conf={"source": "s3://my-bucket/raw/"})
validate() >> trigger_child >> post_trigger_summary()
parent_orchestrator_dag()