Airflow DAG执行失败:Snowflake提示No Active Warehouse
问题
在Airflow中使用Variable.get()方法获取Snowflake的连接ID、角色、数据库、架构及仓库信息时,运行DAG出现「No Active Warehouse」错误;但直接硬编码指定这些信息时,DAG可正常执行。
DAG代码片段
from airflow.contrib.hooks.snowflake_hook import SnowflakeHook from airflow.providers.snowflake.operators.snowflake import SnowflakeOperator from airflow.operators.dummy_operator import DummyOperator from airflow.operators.python_operator import PythonOperator from airflow.sensors.external_task_sensor import ExternalTaskSensor from datetime import datetime, timedelta import pymsteams import logging load_mstr_str = f"call SP_OD_CTL_ETL_EXEC();" default_args = { 'owner': 'MT', 'depends_on_past': False, 'start_date': datetime(2015, 6, 1), 'email_on_failure': False, 'email': ['notify@gmail.com'], 'email_on_failure': True, 'on_failure_callback': send_failure_msg, 'retries': 0 } dag = DAG( dag_id=DAG_ID, catchup=False, default_args=default_args, description="Execute Stored Procs for tables", schedule_interval="30 11 * * *" ) load_mstr_sp = SnowflakeOperator( task_id='load_mstr_sp', dag=dag, snowflake_conn_id=Variable.get("SNOWFLAKE_CONN_ID"), sql=load_mstr_str, warehouse=Variable.get("SNOWFLAKE_WAREHOUSE"), database=Variable.get("SNOWFLAKE_DATABASE"), schema=Variable.get("SNOWFLAKE_SCHEMA"), role=Variable.get("SNOWFLAKE_ROLE"), on_failure_callback=send_failure_msg, )
报错信息
snowflake.connector.errors.ProgrammingError: 000606 (57P03): 01ab6bc6-0502-c477-0030-c80359694e8a: No active warehouse selected in the current session. Select an active warehouse with the 'use warehouse' command.
解决方案
1. 验证Airflow变量的有效性
- 检查Airflow UI中
SNOWFLAKE_WAREHOUSE等变量的值是否准确,注意Snowflake仓库名称区分大小写,必须与实际环境完全一致。 - 添加测试任务打印变量值,确认加载是否正常:
执行该任务,通过日志确认变量值是否正确。def print_vars(): print(f"Warehouse: {Variable.get('SNOWFLAKE_WAREHOUSE')}") print(f"Database: {Variable.get('SNOWFLAKE_DATABASE')}") print_task = PythonOperator( task_id='print_vars', python_callable=print_vars, dag=dag )
2. 改用运行时模板加载变量
Airflow解析DAG文件时会提前执行Variable.get(),可能导致变量未正确初始化。改用Jinja模板在任务运行时动态获取变量:
load_mstr_sp = SnowflakeOperator( task_id='load_mstr_sp', dag=dag, snowflake_conn_id="{{ var.value.SNOWFLAKE_CONN_ID }}", sql=load_mstr_str, warehouse="{{ var.value.SNOWFLAKE_WAREHOUSE }}", database="{{ var.value.SNOWFLAKE_DATABASE }}", schema="{{ var.value.SNOWFLAKE_SCHEMA }}", role="{{ var.value.SNOWFLAKE_ROLE }}", on_failure_callback=send_failure_msg, )
3. 检查Snowflake连接配置
- 确认Airflow中
SNOWFLAKE_CONN_ID对应的连接是否设置了默认仓库,若连接已配置默认仓库,需保证与变量传入的仓库一致,避免冲突。 - 添加测试任务验证连接参数:
运行任务查看日志,确认当前仓库是否正确设置。def test_snowflake_conn(): hook = SnowflakeHook( snowflake_conn_id=Variable.get("SNOWFLAKE_CONN_ID"), warehouse=Variable.get("SNOWFLAKE_WAREHOUSE"), database=Variable.get("SNOWFLAKE_DATABASE"), schema=Variable.get("SNOWFLAKE_SCHEMA"), role=Variable.get("SNOWFLAKE_ROLE") ) conn = hook.get_conn() cursor = conn.cursor() cursor.execute("SELECT CURRENT_WAREHOUSE()") print(f"Current Warehouse: {cursor.fetchone()[0]}") cursor.close() conn.close() test_task = PythonOperator( task_id='test_snowflake_conn', python_callable=test_snowflake_conn, dag=dag )
4. 在SQL中显式指定仓库
若上述方法无效,可在调用存储过程前显式切换仓库:
# 使用变量直接拼接 load_mstr_str = f""" USE WAREHOUSE {Variable.get("SNOWFLAKE_WAREHOUSE")}; call SP_OD_CTL_ETL_EXEC(); """ # 或使用Jinja模板 load_mstr_str = """ USE WAREHOUSE {{ var.value.SNOWFLAKE_WAREHOUSE }}; call SP_OD_CTL_ETL_EXEC(); """
内容的提问来源于stack exchange,提问作者sagyy_sf
相关产品推荐
相关产品推荐

