Spark读取Parquet嵌套结构时遇ClassCastException错误求助
问题描述
我创建了一个包含如下结构的DataFrame:
StructType(List(StructField("shop",StringType), StructField("items",ArrayType(StructType(List(StructField("id",StringType),StructField("qty",StringType),StructField("price",StringType)))))))
对应的数据为:
Seq(Row("shop1",Array(Row("123","34","555"))))
将该DataFrame写入Parquet文件后,其Schema为:
shop: string items: array element: struct id: string qty: string price: string
DataFrame的show()结果为:
shop1 | [{123,34,555}]
当在另一个Spark程序中读取该Parquet文件时,出现以下错误:
java.lang.ClassCastException: optional binary items (UTF8) is not a group
解决方案
这个错误的核心是读取时Spark推断的Schema与Parquet文件实际存储的Schema不匹配——程序误将items识别为普通字符串类型,而实际它是数组嵌套结构体的复杂类型。可按以下方式解决:
1. 显式指定正确Schema读取
读取Parquet时不要依赖自动推断,手动传入创建DataFrame时使用的正确Schema:
import org.apache.spark.sql.types._ // 定义正确的Schema结构 val targetSchema = StructType(List( StructField("shop", StringType), StructField("items", ArrayType(StructType(List( StructField("id", StringType), StructField("qty", StringType), StructField("price", StringType) )))) )) // 用指定Schema读取文件 val df = spark.read.schema(targetSchema).parquet("你的Parquet文件路径")
2. 检查Spark版本兼容性
如果写入和读取使用的Spark版本差异较大,可能存在Parquet格式的兼容性问题。确保两个程序使用的Spark版本一致,或升级到相互兼容的版本。
3. 验证Parquet文件的实际Schema
可以通过以下方式确认文件的实际Schema是否符合预期:
- Spark代码查看:
spark.read.parquet("你的Parquet文件路径").printSchema()
- 命令行工具(parquet-tools)查看:
parquet-tools schema 你的Parquet文件路径/file.parquet
若实际Schema与预期不符,需重新写入DataFrame,确保写入过程中Schema正确无误。
内容的提问来源于stack exchange,提问作者A B
相关产品推荐
相关产品推荐

