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

Spark读取S3分区Parquet文件列类型不一致报错求助

问题描述

我的S3上有如下格式的分区Parquet目录:

database/year=2022/mth=1/id=1/file1.parquet
database/year=2022/mth=12/id=2/file2.parquet
database/year=2022/mth=3/id=3/file3.parquet
database/year=2022/mth=8/id=4/file4.parquet

这些文件中部分列的数据类型不一致,例如file1的impression列是double类型,file2的该列是decimal类型,其他部分列也存在类似问题。

我尝试用以下Scala代码读取数据:

sparkSession.read()
        .option("header", true)
        .option("basePath", basePath)
        .schema(schema)
        .parquet(basePath + input)
        .filter(col("year").equalTo(key).and(col("mth").isin(value.toArray()))
                .and(col("id").isin(idList.toArray())));

但报错:Column: [impressions], Expected: bigint, Found: DOUBLE

我已经试过mergeSchema参数和读取时指定schema,但都没解决问题,求可行的解决办法。

解决方案

1. 分文件读取并手动转换类型

跳过全局schema校验,逐个读取文件后统一转换到目标类型,再合并数据:

// 枚举所有目标文件路径
val filePaths = List(
  s"${basePath}/database/year=2022/mth=1/id=1/file1.parquet",
  s"${basePath}/database/year=2022/mth=12/id=2/file2.parquet",
  s"${basePath}/database/year=2022/mth=3/id=3/file3.parquet",
  s"${basePath}/database/year=2022/mth=8/id=4/file4.parquet"
)

// 定义最终需要的schema
val targetSchema = StructType(Seq(
  StructField("impression", LongType, nullable = true),
  // 补充其他列的类型定义
  StructField("year", StringType, nullable = true),
  StructField("mth", StringType, nullable = true),
  StructField("id", StringType, nullable = true)
))

// 逐个读取文件、转换类型后合并
val finalDF = filePaths.foldLeft(spark.emptyDataFrame) { (acc, path) =>
  val rawDF = spark.read.parquet(path)
  // 按目标schema转换每一列的类型
  val convertedDF = rawDF.select(
    targetSchema.fields.map { field =>
      col(field.name).cast(field.dataType).alias(field.name)
    }: _*
  )
  acc.unionByName(convertedDF)
}

// 应用过滤条件
val filteredDF = finalDF.filter(
  col("year") === key && 
  col("mth").isin(value.toArray: _*) && 
  col("id").isin(idList.toArray: _*)
)

2. 调整Spark配置放宽校验

修改Spark配置,强制开启schema合并并禁用严格的类型校验:

// 配置Spark参数
sparkSession.conf.set("spark.sql.parquet.mergeSchema", "true")
sparkSession.conf.set("spark.sql.parquet.failOnCorruptFile", "false")
sparkSession.conf.set("spark.sql.parquet.enableVectorizedReader", "false")

// 读取数据(不指定schema,先让Spark自动合并)
val mergedDF = sparkSession.read
  .option("basePath", basePath)
  .parquet(basePath + input)

// 手动转换到目标类型
val targetDF = mergedDF.select(
  col("impression").cast(LongType).alias("impression"),
  // 其他列同理转换
  col("year"), col("mth"), col("id")
)

// 应用过滤
targetDF.filter(
  col("year") === key && 
  col("mth").isin(value.toArray: _*) && 
  col("id").isin(idList.toArray: _*)
)

3. 预处理文件统一类型

如果长期存在类型不一致问题,建议一次性预处理所有文件,将类型统一后再读取:

// 获取所有分区的唯一组合
val partitions = spark.read.parquet(basePath + input)
  .select("year", "mth", "id")
  .distinct()
  .collect()

// 遍历每个分区,读取后转换类型并覆盖写回
partitions.foreach { row =>
  val year = row.getAs[String]("year")
  val mth = row.getAs[String]("mth")
  val id = row.getAs[String]("id")
  val partitionPath = s"${basePath}/database/year=${year}/mth=${mth}/id=${id}"
  
  val rawDF = spark.read.parquet(partitionPath)
  // 按目标schema转换所有列
  val convertedDF = rawDF.select(
    targetSchema.fields.map { field =>
      col(field.name).cast(field.dataType).alias(field.name)
    }: _*
  )
  
  // 覆盖写回原分区路径
  convertedDF.write.mode("overwrite").parquet(partitionPath)
}

// 之后就可以正常读取了
sparkSession.read
  .option("basePath", basePath)
  .schema(targetSchema)
  .parquet(basePath + input)
  .filter(...)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 17:18:22