Spark 3.3.0下PySpark DataFrame无法选择结构体字段求助
Dataproc 2.1(Spark 3.3.0)读取Hudi表结构体字段触发ClassCastException的解决方法
问题场景
在GCP Dataproc上运行PySpark任务读取Hudi Parquet表,表结构如下:
root |-- id: string (nullable = true) |-- data: struct (nullable = true) | |-- name: string (nullable = true) | |-- age: string (nullable = true) | |-- city: string (nullable = true) |-- p: integer (nullable = true)
在Dataproc 2.0(Spark 3.1.3)中,以下代码可正常选择结构体字段:
df = df.select( col("id"), col("data.name") )
但迁移到Dataproc 2.1(Spark 3.3.0)后,选择结构体字段会抛出如下错误,仅普通列(如id)可正常使用:
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 516): java.lang.ClassCastException: class org.apache.spark.sql.catalyst.expressions.UnsafeRow cannot be cast to class org.apache.spark.sql.vectorized.ColumnarBatch (org.apache.spark.sql.catalyst.expressions.UnsafeRow and org.apache.spark.sql.vectorized.ColumnarBatch are in unnamed module of loader 'app')
原因分析
Spark 3.3.0默认开启了Parquet向量列读取(Vectorized Reader),该特性优化了原生Parquet文件的读取性能,但适配Spark 3.1的早期Hudi版本,与Spark 3.3的向量读取逻辑存在兼容性问题,导致结构体字段读取时出现类型转换异常。
解决方案
方案1:关闭Parquet向量读取
添加Spark配置禁用向量读取功能,兼容原有Hudi表的读取逻辑:
- 在PySpark代码中配置:
from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("HudiStructFieldRead") \ .config("spark.sql.parquet.enableVectorizedReader", "false") \ .getOrCreate()
- 或在创建Dataproc集群时通过属性指定:
gcloud dataproc clusters create <your-cluster-name> \ --version 2.1 \ --properties spark:spark.sql.parquet.enableVectorizedReader=false
方案2:升级Hudi版本至兼容Spark 3.3的版本
若需保留向量读取的性能优势,可将Hudi版本升级到0.12.0及以上(该版本开始全面适配Spark 3.3.x)。在Dataproc集群中通过指定Hudi依赖包实现升级:
gcloud dataproc clusters create <your-cluster-name> \ --version 2.1 \ --jars gs://hudi-packages/hudi-spark3.3-bundle_2.12-0.12.0.jar
验证
应用上述配置后,重新执行select(col("data.name"))类的代码,即可正常读取结构体字段,不再抛出类型转换异常。
内容的提问来源于stack exchange,提问作者majain
相关产品推荐
相关产品推荐

