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

Delta Lake在PySpark与Python3环境中的行为差异问题排查

GCP Dataproc PySpark读写Delta表报ClassCastException问题

问题场景

已启动GCP Dataproc集群,通过关联Jupyter Notebook及SSH连接的主节点Shell执行操作:

  • Python3 Notebook中可正常读写Delta表
  • PySpark Notebook及PySpark Shell执行完全相同代码时,写入Delta表环节报错

执行代码

from pyspark.sql import SparkSession

# spark.stop

spark = SparkSession.builder \
    .appName('test_session_3') \ # 跨Notebook/Shell运行时已修改会话名称
    .config("spark.jars.packages", "io.delta:delta-storage:2.3.0") \
    .config("spark.jars.packages", "io.delta:delta-core_2.12:2.3.0") \
    .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \
    .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") \
    .getOrCreate()

from pyspark.sql.functions import *
from pyspark.sql.types import StructType,StructField, StringType, IntegerType

data = [(1,"James","","Smith","36636","M",3000),
    (2,"Michael","Rose","","40288","M",4000)
  ]

schema = StructType([ \
    StructField("emp_id",IntegerType(),True), \
    StructField("firstname",StringType(),True), \
    StructField("middlename",StringType(),True), \
    StructField("lastname",StringType(),True), \
    StructField("id", StringType(), True), \
    StructField("gender", StringType(), True), \
    StructField("salary", IntegerType(), True) \
  ])
 
df_employees = spark.createDataFrame(data=data,schema=schema)
df_employees.printSchema()
df_employees.show()

# 写入Delta表时在PySpark环境报错
df_employees.write.format("delta").mode("append").saveAsTable("employee_data_3") 

# 读取Delta表
table_path = '/user/hive/warehouse/employee_data_3'
delta_data_read = spark.read.format("delta").load(table_path)
delta_data_read.count()

错误信息

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

解决方案

1. 修复SparkSession的Jar包配置冲突

重复设置spark.jars.packages会导致后一个配置覆盖前一个,正确做法是将多个依赖包用逗号分隔,写在同一个配置项中:

spark = SparkSession.builder \
    .appName('test_session_3') \
    .config("spark.jars.packages", "io.delta:delta-core_2.12:2.3.0,io.delta:delta-storage:2.3.0") \
    .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \
    .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") \
    .getOrCreate()

2. 验证Spark与Delta版本兼容性

Delta 2.3.0仅适配Spark 3.3.x版本,需确认Dataproc集群的Spark版本:

  • 在Shell执行spark-submit --version查看当前Spark版本
  • 版本不匹配时,要么调整Spark版本,要么更换对应Delta包(例如Spark 3.2.x对应Delta 2.2.x)

3. 确保SparkSession环境干净

如果之前存在未关闭的SparkSession实例,先停止再创建,取消注释spark.stop():

spark.stop()

spark = SparkSession.builder \
    # 配置项...
    .getOrCreate()

4. 切换Delta表写入方式(可选)

若Hive元数据配置存在异常,可跳过saveAsTable,直接将Delta表保存到指定路径:

# 替代saveAsTable的写法
df_employees.write.format("delta").mode("append").save("/user/hive/warehouse/employee_data_3")

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 05:07:44