子 DAG のトリガー
いくつかのワークフローで同じようなコンポーネントが使われていることに気づき、共通のタスクをそれぞれ独立した DAG として切り出せると判断しました。こうすることで、複数のコピーを管理することなく、必要なときに該当コンポーネントを実行できるようになります。現在のワークフロー内のタスクから子 DAG を実行してみましょう。
dag、task、datetime はあらかじめインポートされています。
この演習はコースの一部です
Python で学ぶ Apache Airflow 入門
演習の手順
- ワークフロー内から DAG を起動するために必要なオペレーターをインポートしてください。
child_pipelineという名前の DAG をトリガーするようにオペレーターを設定してください。- 親 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()