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

PySpark读取Parquet遇数据类型不匹配问题求解决方案(无需重写文件)

问题

将一批CSV文件转换为Parquet文件后,加载数据到PySpark DataFrame时遇到类型不兼容错误:

  1. 转换CSV为Parquet的代码:
df = spark.read.options(delimiter='|', inferSchema=True).csv(csv_folder, header=True)
df.write.mode("overwrite").parquet("parquet_folder/")
  1. 读取Parquet文件的代码:
df = spark.read.parquet("parquet_folder/*.parquet")

根据printSchema()输出,字段num_A被识别为integer类型,能成功运行df.show(5),但执行test = df.limit(5)时报错:

Parquet column cannot be converted in file ... Column: [num_A], Expected: int, Found: INT64

若读取时指定schema为LongType:

from pyspark.sql.types import StructType, StructField, LongType

schema = StructType([StructField("num_A", LongType())])
df = spark.read.schema(schema).parquet("parquet_folder/*.parquet")

运行df.show(5)又报错:

Parquet column cannot be converted in file ...(another parquet file) Column: [num_A], Expected: bigint, Found: INT32

核心问题:num_A字段在不同Parquet文件中同时存在INT32(对应PySpark的IntegerType)和INT64(对应LongType)类型,无法通过单一schema读取,如何无需重写Parquet文件解决?

解决方案

方法1:开启mergeSchema自动合并类型

Spark的mergeSchema参数会扫描所有Parquet文件的元数据,自动将同一字段的兼容类型合并为更宽泛的类型(INT32和INT64会合并为INT64/LongType),之后可按需转换为目标类型:

# 开启mergeSchema读取所有文件
df = spark.read.option("mergeSchema", "true").parquet("parquet_folder/")

# 若需要将类型转为IntegerType,需确保数据无超出INT32范围的值
df = df.withColumn("num_A", df["num_A"].cast("integer"))

方法2:通过Spark配置强制兼容类型转换

设置Spark配置项,允许自动处理整数类型的兼容转换:

# 配置读取时自动修正整数类型兼容问题
spark.conf.set("spark.sql.parquet.int96RebaseModeInRead", "CORRECTED")
spark.conf.set("spark.sql.parquet.int64AsTimestamp", "false")

df = spark.read.parquet("parquet_folder/")

方法3:逐个读取文件并统一类型后合并

逐个读取每个Parquet文件,强制转换num_A为统一类型,再合并所有DataFrame:

import os
from pyspark.sql.types import LongType

# 获取所有Parquet文件路径
parquet_files = [os.path.join("parquet_folder", f) for f in os.listdir("parquet_folder") if f.endswith(".parquet")]

# 逐个读取并转换类型
dfs = []
for file in parquet_files:
    temp_df = spark.read.parquet(file)
    temp_df = temp_df.withColumn("num_A", temp_df["num_A"].cast(LongType()))
    dfs.append(temp_df)

# 合并所有DataFrame
final_df = dfs[0]
for df in dfs[1:]:
    final_df = final_df.union(df)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 16:14:55