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

Spark Scala修改Parquet列值后如何保留原始Schema?

保留Parquet文件原始Decimal列Schema的解决方案

问题概述

用这段Scala代码修改Parquet列值后,列类型从原DECIMAL(14,3)变为double,导致查询报错:

var df = spark.read.parquet(sourcePath)
val newDf = df.withColumn("my_num_field", lit(11.10))
newDf.write.parquet(targetPath)

具体问题

  • 原始列Schema(Parquet元数据):
"type" : [ "null", {
      "type" : "fixed",
      "name" : "my_num_field",
      "size" : 16,
      "logicalType" : "decimal",
      "precision" : 14,
      "scale" : 3
    }
  • 写入后新列Schema:
"type" : "double"
  • 查询报错:
incompatible Parquet schema for column 'my_db.my_table.my_num_field'. Column type: DECIMAL(14,3), Parquet schema: required double my_num_field

已尝试无效的方法

  • 在lit(11.10)后追加.cast(DecimalType(14, 3))
  • 读写时设置.option("overwriteSchema","false")
  • 设置mergeSchema为false(搭配或不搭配overwriteSchema)

解决方案

根本原因

lit(11.10)默认生成Double类型值,即便后续cast为Decimal,Spark也会丢失Parquet特有的Fixed-length Decimal元数据(逻辑类型、精度刻度绑定),导致写入时Schema变更。

两种可行实现

方法1:基于原始列Schema构造Decimal值

直接从原始列获取Decimal类型参数,构造匹配的字面量:

import org.apache.spark.sql.types.DecimalType
import org.apache.spark.sql.functions.lit
import org.apache.spark.sql.types.Decimal

val df = spark.read.parquet(sourcePath)
// 提取原始列的Decimal类型信息
val decimalType = df.schema("my_num_field").dataType.asInstanceOf[DecimalType]
// 构造符合精度、刻度的Decimal值
val targetDecimal = Decimal(11.10, decimalType.precision, decimalType.scale)
// 修改列值,保留原始类型
val newDf = df.withColumn("my_num_field", lit(targetDecimal))
// 写入文件
newDf.write.parquet(targetPath)

方法2:使用SQL表达式声明Decimal类型

通过expr直接用SQL语法指定Decimal类型,避免类型转换丢失元数据:

import org.apache.spark.sql.functions.expr

val df = spark.read.parquet(sourcePath)
// 用SQL CAST明确指定DECIMAL(14,3)类型
val newDf = df.withColumn("my_num_field", expr("CAST(11.10 AS DECIMAL(14,3))"))
newDf.write.parquet(targetPath)

验证方式

写入后检查列类型是否匹配:

val resultDf = spark.read.parquet(targetPath)
println(resultDf.schema("my_num_field").dataType)
// 预期输出:DecimalType(14,3)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 01:00:56