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

AWS Glue读取Snowflake正常但写入失败:未选择活跃仓库

解决Spark Glue写入Snowflake时的"No active warehouse selected"错误

已通过Spark Glue成功连接Snowflake并读取数据为DataFrame,但执行写入操作时抛出Py4JJavaError,Snowflake会话历史提示:

No active warehouse selected in the current session. Select an active warehouse with the 'use warehouse' command

代码中已在sfOptions配置了sfWarehouse参数,但写入仍失败。

原因分析

  1. 会话复用问题:代码中调用SparkSession.builder().getOrCreate()创建了新的SparkSession,该会话未继承原会话的Snowflake仓库配置,导致写入操作使用的会话无活跃仓库。
  2. 参数缺失隐患:getResolvedOptions未包含ROLE参数,但sfOptions中引用了args['ROLE'],可能导致配置加载异常。

解决方案

  • 修正enablePushdownSession的调用,使用当前已初始化的SparkSession而非新建会话
  • 确保getResolvedOptions包含所有需要的参数(如ROLE)
  • 验证sfWarehouse参数的拼写正确性(Snowflake对象名称大小写敏感)

修正后的代码示例

import sys
from awsglue.transforms import *
from awsglue.utils import getResolvedOptions
from pyspark.context import SparkContext
from awsglue.context import GlueContext
from awsglue.job import Job
from py4j.java_gateway import java_import

## @params: [JOB_NAME, URL, WAREHOUSE, DB, SCHEMA, USERNAME, PASSWORD, ROLE]
SNOWFLAKE_SOURCE_NAME = "net.snowflake.spark.snowflake"
# 新增ROLE到参数列表,避免KeyError
args = getResolvedOptions(sys.argv, ['JOB_NAME', 'URL', 'WAREHOUSE', 'DB', 'SCHEMA', 'USERNAME', 'PASSWORD', 'ROLE'])

sc = SparkContext()
glueContext = GlueContext(sc)
spark = glueContext.spark_session
job = Job(glueContext)

try:
    job.init(args['JOB_NAME'], args)
except Exception as e:
    pass

java_import(spark._jvm, SNOWFLAKE_SOURCE_NAME)
# 使用当前已有的sparkSession而非新建,确保配置继承
spark._jvm.net.snowflake.spark.snowflake.SnowflakeConnectorUtils.enablePushdownSession(spark._jvm.sparkSession)

sfOptions = {
    "sfURL" : args['URL'],
    "sfUser" : args['USERNAME'],
    "sfPassword" : args['PASSWORD'],
    "sfDatabase" : args['DB'],
    "sfSchema" : args['SCHEMA'],
    "sfWarehouse" : args['WAREHOUSE'],
    "sfRole" : args['ROLE']
}

# 读取数据
df = spark.read.format(SNOWFLAKE_SOURCE_NAME).options(**sfOptions).option("dbtable", "<>").load().select('<>')
print(df.printSchema())
print(df.show())

# 写入数据
df.write.format(SNOWFLAKE_SOURCE_NAME).options(**sfOptions).option("dbtable", "GLUE_DEMO").mode("append").save()

额外验证步骤

  1. 确认Snowflake仓库处于运行状态(未暂停)
  2. 检查当前用户对指定仓库有USAGE权限
  3. 执行代码后查看Snowflake会话历史,确认会话已自动选中目标仓库

内容的提问来源于stack exchange,提问作者Raju

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 22:18:10