Airflow BigQueryInsertJobOperator引用同桶非dags路径SQL文件配置问题
问题根因
Jinja 模板的 {% include %} 指令默认仅检索 Airflow 配置的 DAG 目录下的文件。当使用 GCS 作为 DAG 存储后端时,Airflow 的 GCS DAG 同步逻辑只会将桶内dags/路径加入模板搜索白名单,跨目录的相对路径引用会被路径安全规则拦截,无法访问dags目录外的GCS对象,因此加载失败。
可行实现方案
方案1:扩展DAG模板搜索路径
适用场景:Airflow集群的scheduler、worker节点已通过GCSFuse等方式将整个GCS存储桶挂载到本地文件系统(Google Cloud Composer托管环境默认支持该挂载)。
在DAG初始化时通过template_searchpath参数将SQL脚本所在目录加入Jinja搜索范围,即可直接引用脚本文件,无需写跨目录相对路径,示例代码:
from airflow import DAG from airflow.utils.dates import days_ago from airflow.providers.google.cloud.operators.bigquery import BigQueryInsertJobOperator import os PATH_TO_UPLOAD_FILE_PREFIX = os.environ.get("GCP_GCS_PATH_TO_UPLOAD_FILE_PREFIX", "Test-Processing/") PROJECT_ID = os.environ.get("GCP_PROJECT_ID", "project-id") BUCKET_1 = os.environ.get("GCP_GCS_BUCKET_1", "bucket-name") # 替换为Scripts目录在Airflow运行环境中的本地绝对路径 # Cloud Composer默认挂载路径为 /home/airflow/gcs/Scripts SCRIPTS_LOCAL_PATH = "/home/airflow/gcs/Scripts" with DAG( dag_id='dag_id', default_args=default_args, schedule_interval="@daily", start_date=days_ago(1), catchup=False, template_searchpath=[SCRIPTS_LOCAL_PATH] ) as dag: insert_job_operator = BigQueryInsertJobOperator( task_id='insert_job_operator', configuration={ "query": { "query": "{% include 'script.sql' %}", "useLegacySql": False, } } ) insert_job_operator
注意:未配置GCS本地挂载的环境不要使用该方案,优先选择方案2。
方案2:动态从GCS拉取SQL脚本
适用场景:所有Airflow部署环境,无目录挂载要求,兼容性最强,不受DAG存储路径限制。
不依赖Jinja的include机制,直接通过GCS Hook拉取指定路径的SQL内容,传给BigQuery算子即可,不管SQL文件存在同桶的哪个路径都能正常加载,同时支持SQL内的Jinja动态变量渲染,示例代码:
from airflow import DAG from airflow.utils.dates import days_ago from airflow.providers.google.cloud.operators.bigquery import BigQueryInsertJobOperator from airflow.providers.google.cloud.hooks.gcs import GCSHook import os PATH_TO_UPLOAD_FILE_PREFIX = os.environ.get("GCP_GCS_PATH_TO_UPLOAD_FILE_PREFIX", "Test-Processing/") PROJECT_ID = os.environ.get("GCP_PROJECT_ID", "project-id") BUCKET_1 = os.environ.get("GCP_GCS_BUCKET_1", "bucket-name") with DAG( dag_id='dag_id', default_args=default_args, schedule_interval="@daily", start_date=days_ago(1), catchup=False ) as dag: insert_job_operator = BigQueryInsertJobOperator( task_id='insert_job_operator', configuration={ "query": { "query": "", "useLegacySql": False, } } ) # 任务执行前从GCS拉取SQL,自动渲染Jinja变量 def pull_sql_from_gcs(context): gcs_hook = GCSHook(gcp_conn_id="google_cloud_default") # 替换为SQL文件在GCS上的完整对象路径,无需关联dags目录 sql_content = gcs_hook.download( bucket_name=BUCKET_1, object_name="Scripts/script.sql" ).decode("utf-8") # 渲染SQL中内置的Jinja变量,如执行日期{{ ds }}、自定义变量等 rendered_sql = context["task"].render_template(sql_content, context) context["task"].configuration["query"]["query"] = rendered_sql insert_job_operator.pre_execute = pull_sql_from_gcs insert_job_operator
如果SQL脚本不需要使用动态变量,可以直接在DAG解析阶段拉取SQL内容,代码更简洁:
gcs_hook = GCSHook(gcp_conn_id="google_cloud_default") sql_content = gcs_hook.download(bucket_name=BUCKET_1, object_name="Scripts/script.sql").decode("utf-8") insert_job_operator = BigQueryInsertJobOperator( task_id='insert_job_operator', configuration={ "query": { "query": sql_content, "useLegacySql": False, } } )
内容的提问来源于stack exchange,提问作者pyFlummox
相关产品推荐
相关产品推荐

