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
相关产品推荐
相关产品推荐

