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参数,但写入仍失败。
原因分析
- 会话复用问题:代码中调用
SparkSession.builder().getOrCreate()创建了新的SparkSession,该会话未继承原会话的Snowflake仓库配置,导致写入操作使用的会话无活跃仓库。 - 参数缺失隐患:
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()
额外验证步骤
- 确认Snowflake仓库处于运行状态(未暂停)
- 检查当前用户对指定仓库有
USAGE权限 - 执行代码后查看Snowflake会话历史,确认会话已自动选中目标仓库
内容的提问来源于stack exchange,提问作者Raju
相关产品推荐
相关产品推荐

