如何在Airflow中获取DAG的run_id并将其作为变量传入任务参数
Airflow获取DAG run_id作为变量的实现方法
有两种常用的实现方式,都可以满足你的需求:
- 方法1:在op_kwargs中直接使用Jinja模板引用
Airflow支持在支持模板渲染的参数中直接用{{ run_id }}获取当前运行实例的run_id,你的示例代码可以修改为:
with DAG( dag_id='question', # 其余DAG配置参数 ) as dag: question = PythonOperator( task_id='question_python', python_callable=question_python, op_kwargs={'run_id': '{{ run_id }}'}, )
之后在你的执行函数中直接接收run_id参数即可使用:
def question_python(run_id, **kwargs): # 直接使用传入的run_id即可 print(f"当前运行实例的run_id为:{run_id}") # 你的业务逻辑
- 方法2:从运行上下文直接提取
不需要在op_kwargs中显式传参,在PythonOperator的执行函数中可以直接从上下文参数中提取run_id:
def question_python(**kwargs): run_id = kwargs["run_id"] print(f"当前运行实例的run_id为:{run_id}") with DAG( dag_id='question', # 其余DAG配置参数 ) as dag: question = PythonOperator( task_id='question_python', python_callable=question_python, )
注:如果使用Airflow 1.x版本,需要在PythonOperator参数中添加provide_context=True才能正常获取上下文。除了run_id之外,你还可以用同样的方式获取dag_id、execution_date、task_id等Airflow内置的运行时变量。
内容的提问来源于stack exchange,提问作者Alexander
相关产品推荐
相关产品推荐

