Airflow中Jinja模板二次解析:如何在DAG中设置SQL参数?
解决Airflow中API拉取SQL的模板参数不解析问题
当前问题的核心是:SQLCheckOperator仅会解析sql参数本身的Jinja模板(即{{ ti.xcom_pull(...) }}这部分),但从XCom拉取到的原始SQL内容里的{{ params.table_name }}模板标记不会被自动二次渲染,导致被当作字面量执行。这种场景完全可以实现模板解析,下面是具体解决方案:
方案:通过PythonOperator手动渲染SQL模板
借助Jinja2对拉取到的原始SQL进行二次渲染,替换其中的参数后再传递给SQLCheckOperator执行。
完整代码示例
from jinja2 import Environment from airflow import DAG from airflow.operators.python import PythonOperator from airflow.operators.sql import SQLCheckOperator from airflow.providers.http.operators.http import SimpleHttpOperator from datetime import datetime def render_sql_template(**context): # 从XCom获取API拉取的原始SQL raw_sql = context["ti"].xcom_pull(task_ids="get_sql_query") # 从op_kwargs获取传入的表名参数 table_name = context["table_name"] # 使用Jinja2渲染SQL模板 jinja_env = Environment() rendered_sql = jinja_env.from_string(raw_sql).render(params={"table_name": table_name}) # 将渲染后的SQL推回XCom context["ti"].xcom_push(key="rendered_sql", value=rendered_sql) with DAG( dag_id="sql_api_render_dag", schedule_interval=None, start_date=datetime(2024, 1, 1), catchup=False ) as dag: get_sql_task = SimpleHttpOperator( task_id="get_sql_query", endpoint="path/to/sql_file.sql", method="GET", http_conn_id="some_http_conn_id", log_response=True ) render_task = PythonOperator( task_id="render_sql", python_callable=render_sql_template, provide_context=True, op_kwargs={"table_name": "test_table"} # 这里传入需要替换的表名,也可从DAG参数/上游任务获取 ) sql_check_task = SQLCheckOperator( task_id="check_something", sql="{{ ti.xcom_pull(key='rendered_sql') }}", conn_id="postgres_default" ) get_sql_task >> render_task >> sql_check_task
说明
SimpleHttpOperator拉取原始SQL后,将内容存入XCom;PythonOperator中调用Jinja2的from_string方法加载原始SQL作为模板,传入params参数完成变量替换;- 渲染后的SQL重新存入XCom,供
SQLCheckOperator直接执行,此时SQL中的表名已被正确替换。
如果需要动态传入table_name,可以将其设为DAG的参数(通过dag_run.conf传递),或者从其他上游任务的XCom中获取,只需修改render_sql_template函数的参数来源即可。
内容的提问来源于stack exchange,提问作者sann05
相关产品推荐
相关产品推荐

