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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 00:09:22