如何在Apache Airflow任务中访问@dag装饰器的params参数
在Airflow中访问@dag装饰器的params参数并传递给任务
1. 在DAG定义函数中访问params参数
你可以通过**kwargs中的dag对象直接获取@dag装饰器里定义的params:
from airflow.decorators import dag from datetime import datetime @dag( dag_id='DAG_ID', start_date=datetime(2023, 6, 1), params={ 'id': 'ID_A' } ) def main(**kwargs): # 从dag对象中获取params里的id值 target_id = kwargs['dag'].params['id'] print(target_id) # 输出 'ID_A'
2. 将params参数传递给DAG任务
如果要把这个参数传给具体的Operator任务,有两种常用方式:
方式一:直接在任务的op_kwargs中传递
from airflow.decorators import dag, task from datetime import datetime @dag( dag_id='DAG_ID', start_date=datetime(2023, 6, 1), params={ 'id': 'ID_A' } ) def main(**kwargs): target_id = kwargs['dag'].params['id'] @task def process_id(task_id): print(f"Processing ID: {task_id}") # 将参数传入任务 process_id(task_id=target_id)
方式二:在任务函数中通过上下文(context)访问params
这种方式不需要在定义任务时显式传递,任务可以直接从上下文读取dag的params:
from airflow.decorators import dag, task from datetime import datetime @dag( dag_id='DAG_ID', start_date=datetime(2023, 6, 1), params={ 'id': 'ID_A' } ) def main(): @task def process_id(**context): # 从上下文的params中获取值 target_id = context['params']['id'] print(f"Processing ID: {target_id}") process_id() main()
补充说明
- @dag装饰器中的
params会作为DAG的默认配置参数,当手动触发DAG时,你也可以在触发界面修改这些参数的值,任务会优先使用触发时传入的参数。 - 如果是在传统的Operator(比如PythonOperator)中使用,同样可以通过
provide_context=True开启上下文访问,或者在op_kwargs中传递参数。
内容的提问来源于stack exchange,提问作者Kirill Mukhin
相关产品推荐
相关产品推荐

