PySpark read.csv读取CSV文件格式异常问题求助
解决Spark读取CSV时Schema推断异常的问题
问题原因
Spark的inferSchema机制和pandas逻辑不同:Spark默认采样部分行推断类型,若某列存在混合值(比如多数是数值但夹杂字符串),或者采样没覆盖到所有类型的行,就会出现Schema和预期不符的情况;而pandas的推断逻辑更兼容复杂数据场景。
解决方案
- 手动指定Schema(最可靠)
直接参考pandas读取后的列类型,定义Spark的StructType,彻底避免自动推断的误差。示例代码:
from pyspark.sql.types import StructType, StructField, StringType, IntegerType, FloatType, TimestampType # 对照pandas的Schema定义对应Spark类型,按需补充所有列 custom_schema = StructType([ StructField("CRASH DATE", TimestampType(), nullable=True), StructField("CRASH TIME", StringType(), nullable=True), StructField("BOROUGH", StringType(), nullable=True), StructField("ZIP CODE", StringType(), nullable=True), # ZIP可能含非数值字符,用字符串更稳妥 StructField("LATITUDE", FloatType(), nullable=True), StructField("LONGITUDE", FloatType(), nullable=True), StructField("NUMBER OF PERSONS INJURED", IntegerType(), nullable=True), # 其他列依次添加 ]) data = spark.read.csv('Data/Motor_Vehicle_Collisions_-_Crashes.csv', schema=custom_schema, header=True)
- 调整Spark的采样参数
如果想保留自动推断,可以增大采样比例或行数,让Spark扫描更多数据来准确判断类型:
# 方式1:设置采样比例为1.0(全量扫描,适合数据集不大的情况) data = spark.read.csv('Data/Motor_Vehicle_Collisions_-_Crashes.csv', inferSchema=True, header=True, samplingRatio=1.0) # 方式2:修改全局采样行数参数 spark.conf.set("spark.sql.csv.inferSchema.sampleSize", "100000") # 按需调整数值 data = spark.read.csv('Data/Motor_Vehicle_Collisions_-_Crashes.csv', inferSchema=True, header=True)
- 先读字符串再转换类型
针对存在混合类型的列,先统一以字符串读取,再手动转换并处理异常值:
# 先读取所有列为字符串类型 data = spark.read.csv('Data/Motor_Vehicle_Collisions_-_Crashes.csv', schema="* string", header=True) # 对需要转换的列逐个处理,比如将受伤人数转为整数,异常值设为null from pyspark.sql.functions import col data = data.withColumn("NUMBER OF PERSONS INJURED", col("NUMBER OF PERSONS INJURED").cast(IntegerType()).otherwise(None))
内容的提问来源于stack exchange,提问作者Aditya Jindal
相关产品推荐
相关产品推荐

