如何向Airflow DAG传递参数以在任务外部使用?
问题描述
现有如下Airflow DAG代码:
import datetime as dt from airflow.decorators import dag, task from airflow.operators.dummy import DummyOperator from airflow.models import Variable @dag(start_date=dt.datetime(2022, 1, 1)) def my_dummy_dag(): d0 = DummyOperator(task_id="d0") for i in range(1, 1+int(Variable.get("number", "3"))): d0 >> DummyOperator(task_id="d" + str(i)) my_dummy_dag = my_dummy_dag()
当前通过Variable控制循环生成的任务数量,但希望无需修改Variable或代码,通过触发DAG时的配置或TriggerDagRunOperator等方式传递参数实现需求。
解决方案
当然可以,以下两种方案均可实现你的需求:
方案一:手动触发时传入配置参数
Airflow支持手动触发DAG时传入自定义配置,只需修改目标DAG代码,从dag_run.conf中读取参数替代Variable即可:
修改后的my_dummy_dag代码:
import datetime as dt from airflow.decorators import dag from airflow.operators.dummy import DummyOperator from airflow.models import DagRun @dag(start_date=dt.datetime(2022, 1, 1)) def my_dummy_dag(): def get_task_count(): # 获取当前DAG运行实例的配置 current_dag_run = DagRun.find(dag_id="my_dummy_dag", execution_date=dt.datetime.now())[-1] # 从配置中取number值,默认3 return int(current_dag_run.conf.get("number", 3)) d0 = DummyOperator(task_id="d0") for i in range(1, 1 + get_task_count()): d0 >> DummyOperator(task_id=f"d{i}") my_dummy_dag = my_dummy_dag()
触发操作步骤:
- 在Airflow UI找到该DAG,点击「Trigger DAG」
- 在弹出窗口的「Configuration」输入框中填入JSON配置:
{"number": 5}(数字按需调整) - 确认触发后,DAG会根据传入的
number生成对应数量的任务
方案二:用TriggerDagRunOperator跨DAG传递参数
如果需要通过其他DAG触发目标DAG并传递参数,可使用TriggerDagRunOperator:
触发方DAG代码
import datetime as dt from airflow.decorators import dag from airflow.operators.trigger_dagrun import TriggerDagRunOperator @dag(start_date=dt.datetime(2022, 1, 1), schedule_interval=None) def trigger_my_dummy_dag(): trigger_task = TriggerDagRunOperator( task_id="trigger_my_dummy", trigger_dag_id="my_dummy_dag", conf={"number": 4}, # 传入控制任务数量的参数 wait_for_completion=False ) trigger_my_dummy_dag = trigger_my_dummy_dag()
目标DAG修改
目标DAG的代码与「方案一」中修改后的my_dummy_dag一致即可,触发方DAG运行时会自动将conf中的参数传递给目标DAG,目标DAG据此生成对应数量的任务。
内容的提问来源于stack exchange,提问作者Sonny D
相关产品推荐
相关产品推荐

