如何在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的模板渲染功能,具体实现步骤如下:
- 配置Jinja2环境:指定SQL文件的加载目录,确保能正确读取模板文件
- 渲染SQL模板:传入DAG中定义的自定义宏,将模板内容替换为最终可执行的查询语句
- 执行查询:使用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
相关产品推荐
相关产品推荐

