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

Spark升级至3.3.2后本地PySpark代码报Py4JJavaError: scala/collection/SeqOps缺失

PySpark本地环境适配Python3.9.5+Spark3.3.2启动报错问题

环境信息

  • 运行环境:WSL Ubuntu-22.04
  • Python版本:3.9.5
  • 依赖库版本:
    • py4j: 0.10.9.5
    • pyspark: 3.3.2
  • Spark启动信息:
____              __
     / __/__  ___ _____/ /__
    _\ \/ _ \/ _ `/ __/  '_/
   /___/ .__/\_,_/_/ /_/\_\   version 3.3.2
      /_/
                        
Using Scala version 2.12.15, OpenJDK 64-Bit Server VM, 1.8.0_362

问题描述

代码在Python3.8.10+Spark3.2.1环境、Databricks 10.4 LTS/12.2 LTS集群均可正常运行,但本地环境启动SparkSession时失败。即使注释掉Delta相关配置,仍触发错误。

测试代码

from pyspark import SparkConf
from pyspark.context import SparkContext
from pyspark.sql import SQLContext
from pyspark.sql.session import SparkSession

conf = SparkConf()
# conf.set("spark.sql.shuffle.partitions", "1")
# conf.set("spark.jars.packages", "io.delta:delta-core_2.13:2.2.0")
# conf.set("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension")
# conf.set("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog")

sc = SparkContext.getOrCreate(conf=conf)

spark_session = SparkSession.builder.master("local").appName("Test").config(conf=conf).getOrCreate()

sqlContext = SQLContext(sc, spark_session)

错误日志

Py4JJavaError                             Traceback (most recent call last)
Cell In[2], line 14
      7 # conf.set("spark.sql.shuffle.partitions", "1")
      8 # conf.set("spark.jars.packages", "io.delta:delta-core_2.13:2.2.0")
      9 # conf.set("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension")
     10 # conf.set("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog")
     12 sc = SparkContext.getOrCreate(conf=conf)
---> 14 spark_session = SparkSession.builder.master("local").appName("testTests").config(conf=conf).getOrCreate()
     16 sqlContext = SQLContext(sc, spark_session)

File ~/.pyenv/versions/test9.5_new/lib/python3.9/site-packages/pyspark/sql/session.py:274, in SparkSession.Builder.getOrCreate(self)
    272     session = SparkSession(sc, options=self._options)
    273 else:
--> 274     getattr(
    275         getattr(session._jvm, "SparkSession$"), "MODULE$"
    276     ).applyModifiableSettings(session._jsparkSession, self._options)
    277 return session

File ~/.pyenv/versions/test9.5_new/lib/python3.9/site-packages/py4j/java_gateway.py:1321, in JavaMember.__call__(self, *args)
   1315 command = proto.CALL_COMMAND_NAME +\
   1316     self.command_header +\
   1317     args_command +\
   1318     proto.END_COMMAND_PART
   1320 answer = self.gateway_client.send_command(command)
-> 1321 return_value = get_return_value(
   1322     answer, self.gateway_client, self.target_id, self.name)
   1324 for temp_arg in temp_args:
   1325     temp_arg._detach()

File ~/.pyenv/versions/test9.5_new/lib/python3.9/site-packages/pyspark/sql/utils.py:190, in capture_sql_exception.<locals>.deco(*a, **kw)
    188 def deco(*a: Any, **kw: Any) -> Any:
    189     try:
--> 190         return f(*a, **kw)
    191     except Py4JJavaError as e:
    192         converted = convert_exception(e.java_exception)

File ~/.pyenv/versions/test9.5_new/lib/python3.9/site-packages/py4j/protocol.py:326, in get_return_value(answer, gateway_client, target_id, name)
    324 value = OUTPUT_CONVERTER[type](answer[2:], gateway_client)
    325 if answer[1] == REFERENCE_TYPE:
--> 326     raise Py4JJavaError(
    327         "An error occurred while calling {0}{1}{2}.\n".
    328         format(target_id, ".", name), value)
    329 else:
    330     raise Py4JError(
    331         "An error occurred while calling {0}{1}{2}. Trace:\n{3}\n".
    332         format(target_id, ".", name, value))

Py4JJavaError: An error occurred while calling o58.applyModifiableSettings.
: java.lang.NoClassDefFoundError: scala/collection/SeqOps
    at io.delta.sql.parser.DeltaSqlParser.<init>(DeltaSqlParser.scala:72)
    at io.delta.sql.DeltaSparkSessionExtension.$anonfun$apply$1(DeltaSparkSessionExtension.scala:79)
    at org.apache.spark.sql.SparkSessionExtensions.$anonfun$buildParser$1(SparkSessionExtensions.scala:272)
    at scala.collection.IndexedSeqOptimized.foldLeft(IndexedSeqOptimized.scala:60)
    at scala.collection.IndexedSeqOptimized.foldLeft$(IndexedSeqOptimized.scala:68)
    at scala.collection.mutable.ArrayBuffer.foldLeft(ArrayBuffer.scala:49)
    at org.apache.spark.sql.SparkSessionExtensions.buildParser(SparkSessionExtensions.scala:271)
    at org.apache.spark.sql.internal.BaseSessionStateBuilder.sqlParser$lzycompute(BaseSessionStateBuilder.scala:138)
    at org.apache.spark.sql.internal.BaseSessionStateBuilder.sqlParser(BaseSessionStateBuilder.scala:137)
    at org.apache.spark.sql.internal.BaseSessionStateBuilder.build(BaseSessionStateBuilder.scala:362)
    at org.apache.spark.sql.SparkSession$.org$apache$spark$sql$SparkSession$$instantiateSessionState(SparkSession.scala:1175)
    at org.apache.spark.sql.SparkSession.$anonfun$sessionState$2(SparkSession.scala:162)
    at scala.Option.getOrElse(Option.scala:189)
    at org.apache.spark.sql.SparkSession.sessionState$lzycompute(SparkSession.scala:160)
    at org.apache.spark.sql.SparkSession.sessionState(SparkSession.scala:157)
    at org.apache.spark.sql.SparkSession$.conf$lzycompute$1(SparkSession.scala:1069)
    at org.apache.spark.sql.SparkSession$.conf$1(SparkSession.scala:1069)
    at org.apache.spark.sql.SparkSession$.applyModifiableSettings(SparkSession.scala:1072)
    at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
    at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
    at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
    at java.lang.reflect.Method.invoke(Method.java:498)
    at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244)
    at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:357)
    at py4j.Gateway.invoke(Gateway.java:282)
    at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132)
    at py4j.commands.CallCommand.execute(CallCommand.java:79)
    at py4j.ClientServerConnection.waitForCommands(ClientServerConnection.java:182)
    at py4j.ClientServerConnection.run(ClientServerConnection.java:106)
    at java.lang.Thread.run(Thread.java:750)
Caused by: java.lang.ClassNotFoundException: scala.collection.SeqOps
    at java.net.URLClassLoader.findClass(URLClassLoader.java:387)
    at java.lang.ClassLoader.loadClass(ClassLoader.java:418)
    at java.lang.ClassLoader.loadClass(ClassLoader.java:351)
    ... 30 more

解决方案

核心原因

scala.collection.SeqOps是Scala 2.13新增的类,而当前Spark基于Scala 2.12.15构建。之前配置的Delta包io.delta:delta-core_2.13:2.2.0是Scala 2.13版本,与Spark的Scala版本不兼容。即使注释了配置,Spark仍可能从本地缓存加载不兼容的Delta依赖,导致类加载失败。

修复步骤

  1. 清理Spark缓存:删除本地缓存目录(默认是~/.ivy2/cache或~/.m2/repository下的Delta相关文件),清除残留的不兼容依赖。
  2. 匹配Delta与Spark的Scala版本:若需使用Delta,将Delta包改为适配Scala 2.12的版本,对应Spark 3.3.2的Delta版本为io.delta:delta-core_2.12:2.2.0,修改配置:
    conf.set("spark.jars.packages", "io.delta:delta-core_2.12:2.2.0")
    
  3. 简化SparkSession初始化:无需手动创建SparkContext和SQLContext,SparkSession会自动处理,简化代码:
    from pyspark.sql import SparkSession
    
    spark_session = SparkSession.builder.master("local") \
        .appName("Test") \
        .config("spark.jars.packages", "io.delta:delta-core_2.12:2.2.0") \
        .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \
        .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") \
        .getOrCreate()
    
  4. 重启环境:清理缓存后重启Python解释器或终端,确保新配置生效。

内容的提问来源于stack exchange,提问作者Michał Tołkacz

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 10:09:52