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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 04:10:40