Spark Session依赖验证与Delta JAR加载问题排查
解决Dataproc PySpark Notebook中Delta JAR未加载的问题
一、验证Spark Session是否已加载Delta JAR
方法1:列出已加载的JAR文件
通过SparkContext的listJars()方法查看所有已加载的JAR,筛选Delta相关条目:
# 筛选并打印Delta相关的JAR delta_jars = [jar for jar in spark.sparkContext.listJars() if "delta" in jar.lower()] print("已加载的Delta JAR:", delta_jars)
若输出为空则Delta JAR未加载;若显示对应路径的JAR则已加载。
方法2:执行Delta表操作验证
尝试写入测试Delta表,若抛出ClassNotFoundException说明JAR未加载,反之则正常:
from pyspark.sql.types import StructType, StructField, StringType # 创建测试DataFrame test_df = spark.createDataFrame([("test",)], schema=StructType([StructField("col1", StringType())])) try: # 尝试写入Delta表 test_df.write.format("delta").mode("overwrite").save("/tmp/test_delta_table") print("Delta操作成功,JAR已正常加载") except Exception as e: print("Delta操作失败,JAR未加载:", str(e))
方法3:检查Spark配置
查看spark.jars.packages配置是否包含Delta依赖:
print("当前spark.jars.packages配置:", spark.sparkContext.getConf().get("spark.jars.packages", "未设置"))
若输出未包含io.delta:delta-core_2.12:2.3.0和io.delta:delta-storage:2.3.0,则依赖未正确配置。
二、PySpark环境中补充Delta JAR的方法
方法1:停止已有Session后重新创建(Notebook内解决)
PySpark Notebook通常会自动预创建Spark Session,直接调用builder.getOrCreate()会复用已有Session,导致新配置不生效。需先停止旧Session再重新创建:
from pyspark.sql import SparkSession # 停止已有的Spark Session(忽略不存在的情况) try: spark.stop() except: pass # 重新创建Session,注意将多个依赖包用逗号分隔,避免重复配置spark.jars.packages spark = SparkSession.builder \ .appName('test_session_21') \ .config("spark.jars.packages", "io.delta:delta-core_2.12:2.3.0,io.delta:delta-storage:2.3.0") \ .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \ .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") \ .getOrCreate()
方法2:集群层面全局配置(推荐)
在Dataproc集群创建时,通过初始化动作或修改spark-defaults.conf添加Delta配置,所有Notebook会自动加载:
- 在集群的
spark-defaults.conf中添加以下配置:spark.jars.packages=io.delta:delta-core_2.12:2.3.0,io.delta:delta-storage:2.3.0 spark.sql.extensions=io.delta.sql.DeltaSparkSessionExtension spark.sql.catalog.spark_catalog=org.apache.spark.sql.delta.catalog.DeltaCatalog - 或在创建集群时使用Dataproc初始化动作自动配置Delta依赖。
方法3:手动添加JAR到ClassPath
若无法修改集群配置,可手动下载Delta JAR包上传至集群(如/usr/lib/spark/jars/目录),然后通过配置指定ClassPath:
spark = SparkSession.builder \ .appName('test_session_21') \ .config("spark.driver.extraClassPath", "/usr/lib/spark/jars/delta-core_2.12-2.3.0.jar:/usr/lib/spark/jars/delta-storage-2.3.0.jar") \ .config("spark.executor.extraClassPath", "/usr/lib/spark/jars/delta-core_2.12-2.3.0.jar:/usr/lib/spark/jars/delta-storage-2.3.0.jar") \ .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \ .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") \ .getOrCreate()
内容的提问来源于stack exchange,提问作者user16798185
相关产品推荐
相关产品推荐

