Airflow代码报错:dag_run.conf被识别为字符串,无字典方法
问题原因
你直接把Airflow模板字符串'{{dag_run.conf}}'赋值给了elastic_loader_mode,但在Python函数内部,这个字符串不会被Airflow自动解析渲染,它只是一个普通字符串,自然没有字典的keys()、values()方法,导致执行报错。
解决方案
有两种正确获取dag_run.conf字典值的方式:
方式1:通过函数的context参数获取
修改函数,让它接收Airflow的context参数,从context中直接拿到实际的dag_run.conf字典:
from typing import Any from airflow.models import DagRun def evaluate(elastic_environ: Any, context): # 若dag_run.conf不存在,默认赋值为空字典 elastic_loader_mode = context.get('dag_run').conf or {} airflow_config_expected = {"elastic_mode": "full"} if elastic_loader_mode: # 直接检查键和对应值,避免多余的keys/values调用 if "elastic_mode" in elastic_loader_mode and elastic_loader_mode["elastic_mode"] == "full": print('full') else: raise ValueError(f"Expected config {airflow_config_expected}, but got {elastic_loader_mode}") else: print('delta')
定义PythonOperator时,需开启provide_context=True让函数能拿到Airflow上下文:
PythonOperator( task_id='evaluate_task', python_callable=evaluate, op_kwargs={'elastic_environ': your_environ_value}, provide_context=True, dag=dag )
方式2:通过op_kwargs传递模板变量
在定义PythonOperator时,把dag_run.conf作为参数传递给函数,Airflow会自动解析模板字符串为字典:
from typing import Any def evaluate(elastic_environ: Any, elastic_loader_mode: dict): airflow_config_expected = {"elastic_mode": "full"} if elastic_loader_mode: if "elastic_mode" in elastic_loader_mode and elastic_loader_mode["elastic_mode"] == "full": print('full') else: raise ValueError(f"Expected config {airflow_config_expected}, but got {elastic_loader_mode}") else: print('delta')
对应的PythonOperator定义:
PythonOperator( task_id='evaluate_task', python_callable=evaluate, op_kwargs={ 'elastic_environ': your_environ_value, 'elastic_loader_mode': '{{ dag_run.conf }}' }, dag=dag )
额外优化建议
- 无需用
bool(elastic_loader_mode)判断,直接if elastic_loader_mode:即可,空字典会被视为False - 用
elastic_loader_mode.get("elastic_mode") == "full"替代原判断逻辑,能避免键不存在时抛出KeyError
内容的提问来源于stack exchange,提问作者Nitin Varghese
相关产品推荐
相关产品推荐

