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

升级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依赖方式混乱导致版本不匹配。

修复步骤:

  1. 统一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,避免依赖冲突。
  2. 清理环境缓存:删除Docker容器内的/tmp/metastore_db目录,清理Spark临时文件和Jupyter内核缓存,避免旧元数据干扰。
  3. 验证Scala版本一致性:确认Spark集群所有节点(master、executor)的Scala版本均为2.12.x,无旧版本Spark/Delta jar包遗留。
  4. 临时应急方案:添加Spark序列化器配置改为Kryo,绕过Java序列化的兼容性问题:
    ("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
    

验证

修改配置后重启Spark集群和Jupyter,重新创建Delta表并查询,确认错误消失。

内容的提问来源于stack exchange,提问作者Yaya

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 17:37:10