如何向Airflow TaskFlow风格的DAG与任务传递op_kwargs?
TaskFlow风格DAG传递自定义参数的方法
在TaskFlow风格下,传递自定义参数的方式比旧风格更灵活,主要有以下几种实现方式:
1. 直接调用任务函数时传递参数
TaskFlow装饰的任务函数可以像普通Python函数一样直接传参,直接覆盖默认值,这是最直观的方式:
修改你的DAG代码如下:
from datetime import datetime from airflow.decorators import dag, task from typing import Dict @dag( start_date=datetime.now(), schedule_interval='@once', catchup=False) def example_taskflow_api(): @task() def extract(value=666) -> Dict[str, int]: order_data= {"1001": 301.27, "1002": 433.21, "1003": value} return order_data # 直接传入自定义参数,覆盖默认的666 extract(value=777) my_dag = example_taskflow_api()
2. 通过DAG级别的params配置(支持UI修改)
如果需要在Airflow UI里动态修改参数,可以在DAG定义时添加params字段,然后在任务函数中通过上下文获取:
from datetime import datetime from airflow.decorators import dag, task from typing import Dict @dag( start_date=datetime.now(), schedule_interval='@once', catchup=False, params={"extract_value": 777} # DAG级别的参数配置 ) def example_taskflow_api(): @task() def extract(**kwargs) -> Dict[str, int]: # 从上下文的params中获取参数 value = kwargs['params']['extract_value'] order_data= {"1001": 301.27, "1002": 433.21, "1003": value} return order_data extract() my_dag = example_taskflow_api()
运行DAG时,你可以在UI的"Run DAG"页面修改extract_value的值,无需改动代码。
3. 任务级别的params配置
如果只想给单个任务配置可修改的参数,可以直接在@task装饰器中指定params:
from datetime import datetime from airflow.decorators import dag, task from typing import Dict @dag( start_date=datetime.now(), schedule_interval='@once', catchup=False) def example_taskflow_api(): # 给当前任务单独配置参数 @task(params={"value": 777}) def extract(**kwargs) -> Dict[str, int]: value = kwargs['params']['value'] order_data= {"1001": 301.27, "1002": 433.21, "1003": value} return order_data extract() my_dag = example_taskflow_api()
这种方式的参数仅对该任务生效,同样支持在UI运行时修改。
内容的提问来源于stack exchange,提问作者BenP
相关产品推荐
相关产品推荐

