Airflow中TemplateNotFound问题:子目录DAG无法读取ConfigSpark.yaml
目录结构
. ├── ConfigSpark.yaml ├── project1 │ ├── dags │ │ └── dag_1.py │ └── sparkjob │ └── spark_1.py └── sparkutils
代码实现
我在dag_1.py中使用SparkKubernetesOperator导入ConfigSpark.yaml,代码如下:
job= SparkKubernetesOperator( task_id = 'job', params=dict( app_name='job', mainApplicationFile='/opt/airflow/dags/project1/sparkjob/spark_1.py', driverCores=1, driverCoreRequest='250m', driverCoreLimit='500m', driverMemory='2G', executorInstances=1, executorCores=2, executorCoreRequest='1000m', executorCoreLimit='1000m', executorMemory='2G' ), application_file='/opt/airflow/dags/ConfigSpark.yaml', kubernetes_conn_id='conn_prd_eks', do_xcom_push=True )
报错信息
运行DAG时返回错误:
jinja2.exceptions.TemplateNotFound: /opt/airflow/dags/ConfigSpark.yaml
疑问
当DAG与ConfigSpark.yaml在同一目录时任务可正常运行,但将DAG放在子目录时就无法运行。已检查values.yaml文件,其中airflowHome为/opt/airflow,defaultAirflowRepository为apache/airflow,请问问题出在哪里?
核心原因
Airflow处理SparkKubernetesOperator的application_file参数时,会将其当作Jinja2模板文件解析,而Jinja2默认的模板查找路径是DAG文件所在的目录,并非你指定的绝对路径。
当DAG和ConfigSpark.yaml同目录时,Jinja2能在当前DAG目录下找到文件;但DAG放到子目录后,Jinja2只会在project1/dags/下查找,找不到上层目录的ConfigSpark.yaml,因此抛出TemplateNotFound错误。
解决方案
有三种可行的解决方式:
修改Jinja2模板搜索路径
在DAG文件开头,修改Airflow的Jinja2环境,添加ConfigSpark.yaml所在目录到模板搜索路径:from jinja2 import FileSystemLoader from airflow.utils.session import create_session from airflow.models import Variable # 获取Airflow的Jinja环境并添加新的模板路径 with create_session() as session: jinja_env = Variable.get_jinja_environment(session=session) jinja_env.loader.searchpath.append('/opt/airflow/dags/')之后再定义
SparkKubernetesOperator,即可正确找到目标文件。调整文件位置或创建软链接
- 直接将
ConfigSpark.yaml复制到project1/dags/目录下,使用相对路径ConfigSpark.yaml即可; - Linux环境下,在
project1/dags/目录下创建指向目标文件的软链接:ln -s /opt/airflow/dags/ConfigSpark.yaml /opt/airflow/dags/project1/dags/ConfigSpark.yaml
- 直接将
直接读取文件内容传入参数
跳过Jinja2模板加载逻辑,直接读取ConfigSpark.yaml的内容,传递给application_file参数(该参数支持传入字符串格式的配置内容):import yaml # 读取配置文件内容 with open('/opt/airflow/dags/ConfigSpark.yaml', 'r') as f: spark_config_content = yaml.safe_load(f) job= SparkKubernetesOperator( task_id = 'job', params=dict( app_name='job', mainApplicationFile='/opt/airflow/dags/project1/sparkjob/spark_1.py', driverCores=1, driverCoreRequest='250m', driverCoreLimit='500m', driverMemory='2G', executorInstances=1, executorCores=2, executorCoreRequest='1000m', executorCoreLimit='1000m', executorMemory='2G' ), application_file=spark_config_content, kubernetes_conn_id='conn_prd_eks', do_xcom_push=True )
内容的提问来源于stack exchange,提问作者OdiumPura

