Airflow中XCom Pull用f-string触发Jinja未定义错误排查
问题原因与修复方案
错误根源
你构建的f-string在生成Jinja表达式时,没有给task_ids的参数值添加引号。当f-string渲染后,最终的Jinja代码变成了:
{{ ti.xcom_pull(task_ids=get_config_for_dag_bla_bla_sbx , key='return_value') }}
Jinja会把get_config_for_dag_bla_bla_sbx识别为一个变量而非字符串常量,但这个变量在Jinja上下文里不存在,因此抛出UndefinedError。
修复代码
只需要在f-string中给previous_task_id包裹单引号,让Jinja能正确识别它是字符串参数:
previous_task_id = 'get_config_for_dag_' + external_dag_id # 给{previous_task_id}添加单引号 conf=f"{{{{ ti.xcom_pull(task_ids='{previous_task_id}' , key='return_value') }}}}"
渲染后的正确Jinja表达式为:
{{ ti.xcom_pull(task_ids='get_config_for_dag_bla_bla_sbx' , key='return_value') }}
完整修正后的循环代码片段
dag_info = ['bla_bla_sbx'] for external_dag_id in dag_info: get_config_for_dag = PythonOperator( task_id="get_config_for_dag_" + external_dag_id, python_callable=get_dict_value, op_kwargs={ 'config_dict': "{{ ti.xcom_pull(task_ids='build_dags_configuration', key='return_value') }}", 'current_dag_id': external_dag_id } ) previous_task_id = 'get_config_for_dag_' + external_dag_id run_dag = TriggerDagRunOperator( task_id="run_dag_" + external_dag_id, wait_for_completion=True, trigger_dag_id=external_dag_id, pool="dag_pool", # 修复此处的f-string conf=f"{{{{ ti.xcom_pull(task_ids='{previous_task_id}' , key='return_value') }}}}", )
内容的提问来源于stack exchange,提问作者Fernando Garcia Dorador
相关产品推荐
相关产品推荐

