Запуск дочернего DAG
Вы замечаете, что некоторые из ваших рабочих процессов используют схожие компоненты, и понимаете, что общие задачи можно вынести в отдельный DAG. Это позволит запускать их по мере необходимости, не поддерживая несколько копий. Вы решаете запустить дочерний DAG из задачи в текущем рабочем процессе.
Компоненты dag, task и datetime уже импортированы за вас.
Это упражнение является частью курса
Введение в Apache Airflow на Python
Инструкции к упражнению
- Импортируйте оператор, необходимый для запуска DAG из текущего рабочего процесса.
- Настройте оператор так, чтобы он запускал DAG с именем
child_pipeline. - Убедитесь, что родительский DAG ожидает завершения запущенного DAG перед продолжением работы.
- Задайте интервал, с которым оператор проверяет, завершил ли работу дочерний DAG.
Интерактивное практическое упражнение
Попробуйте выполнить это упражнение, дополнив этот пример кода.
# 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()