使用get_current_context()执行PythonOperator遇template_fields不存在KeyError求助
问题
使用get_current_context()并通过PythonOperator执行任务时,遇到错误:Variable template_fields does not exist。
代码示例
@task def varfile(regularvalue,previousvalue,dag_instance, **kwargs): if regularvalue: context = get_current_context() varvalue = PythonOperator( task_id=f"varvalue", python_callable=def_varvalue, provide_context= True, op_kwargs = {"dag_name":dag_name,"regularvalue":regularvalue,"previousvalue":int(previousvalue)}, dag=dag_instance ) varvalue.execute(context) varvalue finalval = varfile("{{ params.regvalue }}","{{ params.regvalue_previous_date }}", dag_instance=dag)
报错信息
KeyError: 'Variable template_fields does not exist' File "/home/dag.py", line 100, in varfile varvalue.execute(context) File "/home/airflow/operators/python.py", line 148, in execute self.op_kwargs = determine_kwargs(self.python_callable, self.op_args, context) File "/home/airflow/models/baseoperator.py", line 742, in __setattr__ self.set_xcomargs_dependencies() File "/home/airflow/models/baseoperator.py", line 864, in set_xcomargs_dependencies apply_set_upstream(arg) File "/home/airflow/models/baseoperator.py", line 856, in apply_set_upstream apply_set_upstream(elem) File "/home/airflow/models/baseoperator.py", line 856, in apply_set_upstream apply_set_upstream(elem) File "/home/airflow/models/baseoperator.py", line 857, in apply_set_upstream elif hasattr(arg, "template_fields"): File "/home/airflow/models/taskinstance.py", line 1663, in __getattr__ self.var = Variable.get(item, deserialize_json=True) File "/home/airflow/models/variable.py", line 140, in get raise KeyError(f'Variable {key} does not exist') KeyError: 'Variable template_fields does not exist'
解决方案
核心问题
你在@task装饰的函数内部实例化并执行PythonOperator的做法是错误的。Airflow的@task装饰器已经将函数包装成TaskOperator,内部再创建Operator会触发Airflow的内部属性检查逻辑,当前上下文的TaskInstance对象会误将template_fields当作Airflow Variable去查找,从而抛出错误。
修改步骤
- 直接调用目标函数:无需在任务函数内部创建PythonOperator,直接调用
def_varvalue即可。 - 移除不必要的参数:
@task装饰的函数不需要手动传入dag_instance,Airflow会自动处理任务与DAG的关联。 - 弃用
provide_context=True:Airflow 2.x中该参数已被弃用,推荐通过get_current_context()或**kwargs获取上下文参数。
修改后的代码
@task def varfile(regularvalue, previousvalue, **kwargs): if regularvalue: # 获取上下文(如果def_varvalue需要) context = get_current_context() # 直接调用目标函数,传入所需参数 def_varvalue( dag_name=dag_name, regularvalue=regularvalue, previousvalue=int(previousvalue), **context # 若def_varvalue依赖上下文参数,直接传入 ) # 调用任务,无需传入dag_instance finalval = varfile("{{ params.regvalue }}", "{{ params.regvalue_previous_date }}")
简化版(若无需上下文参数)
如果def_varvalue不需要额外的上下文参数,可进一步简化:
@task def varfile(regularvalue, previousvalue): if regularvalue: def_varvalue( dag_name=dag_name, regularvalue=regularvalue, previousvalue=int(previousvalue) )
内容的提问来源于stack exchange,提问作者Arya
相关产品推荐
相关产品推荐

