Spark写入Delta格式报找不到spark_catalog插件类错误如何解决
运行环境
- Spark版本:3.2.1
- Delta版本:1.2.1(曾尝试使用2.0版本)
问题现象
运行Delta格式写入的入门测试代码时抛出异常,测试代码如下:
from pyspark.sql import SparkSession from delta import * builder = SparkSession.builder.appName("MyApp") \ .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \ .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") spark = configure_spark_with_delta_pip(builder).getOrCreate() data = spark.range(0, 5) data.write.format("delta").save("/tmp/delta-table")
报错信息
- 错误类型:
Py4JJavaError - 核心报错:
org.apache.spark.SparkException: Cannot find catalog plugin class for catalog 'spark_catalog'
即无法找到标识为spark_catalog的catalog对应的插件类。
问题成因
- SparkSession配置未生效:运行环境中存在提前初始化的默认SparkSession实例,
getOrCreate()方法直接返回未配置Delta参数的已有实例,没有加载Delta的catalog相关配置。 - 依赖版本不匹配:使用的Delta包和Spark版本、Scala编译版本不对应,或者环境中存在多个版本的Delta依赖包,导致类加载时无法找到对应
DeltaCatalog类。 - 依赖分发异常:集群模式提交任务时,Delta依赖包没有被正确分发到所有计算节点,节点上缺失对应类。
解决方案
- 校验版本适配关系
Spark 3.2.x 仅兼容Delta 1.2.x、2.0.x版本,且依赖包的Scala版本需和Spark编译使用的Scala版本保持一致(Spark 3.2.x默认使用Scala 2.12编译),禁止跨大版本混用依赖。 - 清理预初始化的Spark实例
如果在Jupyter Notebook、PySpark Shell等预置Spark环境中运行代码,先执行如下命令停止已存在的Spark实例,避免配置被跳过:spark.stop() - 显式指定依赖版本构建SparkSession
调用configure_spark_with_delta_pip时显式传入匹配版本的Delta依赖,避免自动拉取到版本不兼容的包,示例代码:spark = configure_spark_with_delta_pip( builder, extra_packages=["io.delta:delta-core_2.12:1.2.1"] ).getOrCreate() - 排查本地依赖冲突
检查$SPARK_HOME/jars目录下是否存在其他版本的delta-core jar包,如果存在则全部移除,避免类加载优先级冲突导致找不到正确类。 - 验证配置加载结果
SparkSession启动后执行如下命令校验配置是否正确生效:
正常情况下输出应为print(spark.conf.get("spark.sql.catalog.spark_catalog"))org.apache.spark.sql.delta.catalog.DeltaCatalog,如果输出不符合预期,需要排查是否有其他配置文件覆盖了自定义参数。
内容的提问来源于stack exchange,提问作者Mohan
相关产品推荐
相关产品推荐

