PySpark读取Parquet遇数据类型不匹配问题求解决方案(无需重写文件)
问题
将一批CSV文件转换为Parquet文件后,加载数据到PySpark DataFrame时遇到类型不兼容错误:
- 转换CSV为Parquet的代码:
df = spark.read.options(delimiter='|', inferSchema=True).csv(csv_folder, header=True) df.write.mode("overwrite").parquet("parquet_folder/")
- 读取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

