PySpark读取多Parquet文件同列多类型致Schema合并报错处理
Spark读取Parquet同名列类型不兼容兼容方案
问题背景
- 现有跨度4年的大量订单类Parquet文件,原有批量读取逻辑如下:
df = spark.read.option( 'mergeSchema', True).parquet(*list_order).select( 'at', 'order_id', 'items')
- 所有文件Schema一致时逻辑可正常运行,新版接入数据中
quantity字段类型从String变为Float/Double,触发Schema合并失败,报错如下:
Caused by: org.apache.spark.SparkException: Failed to merge fields 'quantity' and 'quantity'. Failed to merge incompatible data types string and double
- 落地约束:不得重生成全量历史Parquet文件,全量重刷4年数据生产环境耗时过高,无落地可行性。
可落地方案(无需修改历史文件)
方案1:读取时显式指定统一Schema(成本最低)
Spark自动Schema合并遇到不兼容类型会直接抛错,提前定义好最终统一的目标Schema传入读取逻辑,Spark会按照指定Schema做字段类型适配,跳过自动合并的类型校验环节。
如果历史数据中quantity字段存储的字符串均为合法数值格式(如"2"、"3.5"),直接指定为DoubleType即可自动完成转换,无额外清洗成本。
示例代码:
from pyspark.sql.types import * # 按实际业务字段补全全量Schema,以下为参考示例 target_schema = StructType([ StructField("at", TimestampType(), True), StructField("order_id", StringType(), True), # 若quantity为items嵌套结构内的字段,需同步定义嵌套层的字段类型 StructField("items", ArrayType(StructType([ StructField("sku_id", StringType(), True), StructField("quantity", DoubleType(), True), # 统一指定为数值类型 # 补充items下其余业务字段 ])), True) ]) # 关闭自动mergeSchema,传入预定义Schema读取 df = spark.read \ .schema(target_schema) \ .parquet(*list_order) \ .select('at', 'order_id', 'items')
补充:如果历史数据存在无法转换为数值的异常字符串(如空串、特殊字符),读取完成后增加一层清洗规则即可,不会影响整体读取流程。
方案2:按Schema版本拆分读取,对齐类型后合并
字段类型变更存在明确的上线时间节点,可按节点将文件列表拆分为历史旧Schema文件、新Schema文件两部分,分别读取后对旧数据做类型转换,字段完全对齐后再做合并,全程不会触发Schema合并冲突。
该方案容错性更高,遇到脏数据可单独针对旧数据做定制化清洗,不会出现全量读取失败的问题。
示例代码:
from pyspark.sql.functions import col, expr # 按时间节点拆分文件路径列表 # old_list_order:类型变更前的所有历史Parquet路径 # new_list_order:类型变更后的新数据Parquet路径 df_old = spark.read.parquet(*old_list_order) # 旧数据类型转换:将String类型的quantity转为Double,嵌套结构用transform处理数组字段 df_old_processed = df_old.select( col("at"), col("order_id"), # 非嵌套场景直接cast即可:col("quantity").cast(DoubleType()).alias("quantity") expr(""" transform(items, item -> named_struct( 'sku_id', item.sku_id, 'quantity', cast(item.quantity as double), 'item_name', item.item_name )) as items """) ) df_new = spark.read.parquet(*new_list_order) # 对齐新数据字段顺序,和处理后的旧数据保持完全一致 df_new_processed = df_new.select('at', 'order_id', 'items') # 按字段名合并两份数据,无类型冲突 final_df = df_old_processed.unionByName(df_new_processed)
风险提示:不要尝试直接修改Parquet文件元数据的方式统一字段类型,该操作直接触碰底层存储文件,极易造成文件损坏、数据不可读的生产事故,稳定性远低于上述两种读取时处理的方案。
内容的提问来源于stack exchange,提问作者AdriHein
相关产品推荐
相关产品推荐

