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依赖,导致类加载失败。
修复步骤
- 清理Spark缓存:删除本地缓存目录(默认是
~/.ivy2/cache或~/.m2/repository下的Delta相关文件),清除残留的不兼容依赖。 - 匹配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") - 简化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() - 重启环境:清理缓存后重启Python解释器或终端,确保新配置生效。
内容的提问来源于stack exchange,提问作者Michał Tołkacz
相关产品推荐
相关产品推荐

