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
相关产品推荐
相关产品推荐

