Airflow中SnowflakeOperator如何使用DAG运行配置传入的schema参数
根因说明
旧版本SnowflakeOperator的schema参数默认不在Jinja模板渲染字段列表中,代码中填写的{{ dag_run.conf['database_schema'] }}不会被Airflow解析为动态模板,而是会被直接作为字符串传递给Snowflake连接,自然找不到对应schema下的存储过程。而PythonOperator的op_kwargs默认支持模板渲染,所以可以正常拿到传入的schema值。
解决方案
方案1:自定义扩展SnowflakeOperator(最稳妥,无兼容性问题)
继承原生SnowflakeOperator,将schema加入模板字段列表即可:
from airflow.providers.snowflake.operators.snowflake import SnowflakeOperator class DynamicSchemaSnowflakeOperator(SnowflakeOperator): # 追加schema到模板渲染字段 template_fields = (*SnowflakeOperator.template_fields, 'schema') # 后续使用自定义Operator即可 test_sp = DynamicSchemaSnowflakeOperator( task_id="test_sp", sql="CALL usp_check_for_null_op()", snowflake_conn_id=SNOWFLAKE_CONN, schema="{{ dag_run.conf['database_schema'] }}" )
方案2:直接在SQL语句中指定schema
SnowflakeOperator的sql参数默认支持模板渲染,你可以直接把动态schema拼到SQL里,不需要改Operator:
test_sp = SnowflakeOperator( task_id="test_sp", sql="CALL {{ dag_run.conf['database_schema'] }}.usp_check_for_null_op()", snowflake_conn_id=SNOWFLAKE_CONN )
方案3:升级Snowflake Provider包
如果你的AWS托管Airflow环境允许升级依赖,将apache-airflow-providers-snowflake升级到2.0.0及以上版本,该版本已经官方将schema加入了SnowflakeOperator的模板字段列表,原生支持动态schema参数。
验证方法
可以在任务详情页的渲染模板(Rendered Template) tab中查看参数实际渲染结果,确认schema是否被替换为传入的配置值,即可快速排查模板渲染问题。
内容的提问来源于stack exchange,提问作者YogeshR
相关产品推荐
相关产品推荐

