Zacznij terazZacznij za darmo

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

Zobacz kurs

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()
Edytuj i uruchom kod