You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Airflow的常规函数或PythonOperator中解析自定义宏

问题描述

我们在GCP项目中使用托管式Airflow。此前使用BigQueryInsertJobOperator执行查询文件时,系统会自动将文件中的user_defined_macros替换为设定值,示例代码如下:

from airflow import DAG
from datetime import datetime
from airflow.providers.google.cloud.operators.bigquery import BigQueryInsertJobOperator

with DAG(
    'test',
    schedule_interval = None,
    start_date = datetime(2022, 1, 1),
    user_defined_macros = {
        "MY_MACRO": "Hello World"
    }
) as dag:

    BigQueryInsertJobOperator(
        task_id = "my_task",
        configuration = {
                            "query": {
                                "query": "{% include '/queries/my_query.sql' %}",
                                "useLegacySql": False,
                            },
                        },
        dag = dag,
    )

因某些原因,我们现在切换为使用常规函数或PythonOperator通过BigQuery客户端执行这些查询,但无法实现自定义宏的解析。以下是目前的代码(无法正常运行):

from airflow import DAG
from datetime import datetime
from google.cloud import bigquery
from airflow.decorators import task

with DAG(
    'test',
    schedule_interval = None,
    start_date = datetime(2022, 1, 1),
    user_defined_macros = {
        "MY_MACRO": "Hello World"
    }
) as dag:
    
    @task
    def test():
        query = open('/home/airflow/gcs/dags/queries/my_query.sql').read()
        bq_client = bigquery.Client()
        bq_client.query(query).result()

    test()
解决方案

要在Python任务中解析Airflow的Jinja2宏,需要手动调用Jinja2的模板渲染功能,具体实现步骤如下:

  1. 配置Jinja2环境:指定SQL文件的加载目录,确保能正确读取模板文件
  2. 渲染SQL模板:传入DAG中定义的自定义宏,将模板内容替换为最终可执行的查询语句
  3. 执行查询:使用BigQuery客户端执行渲染后的SQL

修改后的可运行代码示例:

from airflow import DAG
from datetime import datetime
from google.cloud import bigquery
from airflow.decorators import task
from jinja2 import Environment, FileSystemLoader

with DAG(
    'test',
    schedule_interval=None,
    start_date=datetime(2022, 1, 1),
    user_defined_macros={
        "MY_MACRO": "Hello World"
    }
) as dag:
    
    @task
    def test():
        # 初始化Jinja2环境,指向SQL文件所在目录
        template_env = Environment(loader=FileSystemLoader('/home/airflow/gcs/dags/queries'))
        # 加载模板并传入自定义宏完成渲染
        query = template_env.get_template('my_query.sql').render(**dag.user_defined_macros)
        # 执行BigQuery查询
        bq_client = bigquery.Client()
        bq_client.query(query).result()

    test()

补充说明

如果SQL模板中还用到了Airflow内置宏(比如{{ ds }}),需要传入Airflow上下文变量,可通过provide_context=True参数获取并合并到渲染参数中:

@task(provide_context=True)
def test(**context):
    template_env = Environment(loader=FileSystemLoader('/home/airflow/gcs/dags/queries'))
    # 合并自定义宏与Airflow上下文变量
    render_vars = {**dag.user_defined_macros, **context}
    query = template_env.get_template('my_query.sql').render(**render_vars)
    bq_client = bigquery.Client()
    bq_client.query(query).result()

注意:托管式Airflow中需确保SQL文件路径正确,可根据实际目录调整FileSystemLoader的参数。

内容的提问来源于stack exchange,提问作者Shahid Thaika

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.13 09:20:42