升级Spark 3.5.0与Delta 3.1.0后Delta表操作失败问题排查
Spark 3.5.0 + Delta 3.1.0 升级后查询Delta表报错问题
环境与问题背景
- 原正常运行环境:
- Ubuntu 20.04 (WSL)
- openjdk:17.0.2
- Scala 2.12
- Spark 3.4.0
- Spark-delta 2.4.0
- JupyterLab
- 升级配置:Spark 3.5.0 + 官方适配的Delta 3.1.0版本
- 问题:升级后创建或查询Delta表时抛出
ClassCastException,核心错误为cannot assign instance of java.lang.invoke.SerializedLambda to field org.apache.spark.sql.catalyst.expressions.ScalaUDF.f of type scala.Function1
Spark Session 创建代码
spark_conf.setAll( [ ("spark.master", "spark://spark-master:7077"), ("spark.app.name", "spark_app"), ("spark.driver.memory", "4g"), ("spark.submit.deployMode", "client"), ("spark.ui.showConsoleProgress", "true"), ("spark.eventLog.enabled", "false"), ("spark.logConf", "false"), ( "spark.jars", "/usr/lib/delta-core_2.12-3.1.0.jar", ), ("spark.driver.extraJavaOptions", "-Djava.net.useSystemProxies=true"), ("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension"), ( "spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog", ), ( "javax.jdo.option.ConnectionURL", f"jdbc:derby:;databaseName=/tmp/metastore_db;create=true", ), ("spark.sql.catalogImplementation", "hive"), ] ) builder = SparkSession.builder.config(conf=spark_conf) spark_session = configure_spark_with_delta_pip(builder).getOrCreate()
Delta表查询代码
df = spark_session.sql(f"""select * from delta_table;""") df.show()
报错信息
Py4JJavaError: An error occurred while calling o57.sql. : org.apache.spark.SparkException: Job aborted due to stage failure: Task 0 in stage 1.0 failed 4 times, most recent failure: Lost task 0.3 in stage 1.0 (TID 4) (172.20.0.4 executor 0): java.lang.ClassCastException: cannot assign instance of java.lang.invoke.SerializedLambda to field org.apache.spark.sql.catalyst.expressions.ScalaUDF.f of type scala.Function1 in instance of org.apache.spark.sql.catalyst.expressions.ScalaUDF at java.base/java.io.ObjectStreamClass$FieldReflector.setObjFieldValues(ObjectStreamClass.java:2227) at java.base/java.io.ObjectStreamClass$FieldReflector.checkObjectFieldValueTypes(ObjectStreamClass.java:2191) at java.base/java.io.ObjectStreamClass.checkObjFieldValueTypes(ObjectStreamClass.java:1478) at java.base/java.io.ObjectInputStream$FieldValues.defaultCheckFieldValues(ObjectInputStream.java:2690) at java.base/java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2497) at snipped .... at java.base/java.io.ObjectInputStream.readObject(ObjectInputStream.java:472) at scala.collection.immutable.List$SerializationProxy.readObject(List.scala:527) at jdk.internal.reflect.GeneratedMethodAccessor4.invoke(Unknown Source) at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.base/java.lang.reflect.Method.invoke(Method.java:568) at snipped .. java.base/java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1744) at java.base/java.io.ObjectInputStream.readObject(ObjectInputStream.java:514) at java.base/java.io.ObjectInputStream.readObject(ObjectInputStream.java:472) at org.apache.spark.serializer.JavaDeserializationStream.readObject(JavaSerializer.scala:87) at org.apache.spark.serializer.JavaSerializerInstance.deserialize(JavaSerializer.scala:129) at org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:90) at org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:54) at org.apache.spark.TaskContext.runTaskWithListeners(TaskContext.scala:161) at org.apache.spark.scheduler.Task.run(Task.scala:141) at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$4(Executor.scala:620) at org.apache.spark.util.SparkErrorUtils.tryWithSafeFinally(SparkErrorUtils.scala:64) at org.apache.spark.util.SparkErrorUtils.tryWithSafeFinally$(SparkErrorUtils.scala:61) at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:94) at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:623) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635) at java.base/java.lang.Thread.run(Thread.java:833)
问题分析与解决
这是序列化兼容性问题,核心原因是Spark 3.5.0调整了Scala UDF的序列化逻辑,同时当前环境中Delta依赖方式混乱导致版本不匹配。
修复步骤:
- 统一Delta依赖方式:不要同时使用
spark.jars指定本地jar包和configure_spark_with_delta_pip的pip安装方式,二选一:- 若用pip安装:删除
spark.jars配置项,确保安装的delta-spark版本为3.1.0(与Spark 3.5.0严格适配)。 - 若用本地jar包:同时引入
delta-core_2.12-3.1.0.jar和delta-storage_2.12-3.1.0.jar,且不要调用configure_spark_with_delta_pip,避免依赖冲突。
- 若用pip安装:删除
- 清理环境缓存:删除Docker容器内的
/tmp/metastore_db目录,清理Spark临时文件和Jupyter内核缓存,避免旧元数据干扰。 - 验证Scala版本一致性:确认Spark集群所有节点(master、executor)的Scala版本均为2.12.x,无旧版本Spark/Delta jar包遗留。
- 临时应急方案:添加Spark序列化器配置改为Kryo,绕过Java序列化的兼容性问题:
("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
验证
修改配置后重启Spark集群和Jupyter,重新创建Delta表并查询,确认错误消失。
内容的提问来源于stack exchange,提问作者Yaya
相关产品推荐
相关产品推荐

