PySpark重分区S3中Parquet文件失败及类型问题求助
问题根源与解决方案
这个错误的核心原因是你的col1列存在混合数据类型(部分行是浮点型数值,部分是字符串型的数字),而Parquet的字典编码机制在处理这种类型不一致的列时,会出现解码失败的情况——尤其是在执行分区、shuffle这类需要对列数据进行编码/解码的操作时,Spark无法处理同一列里的混合类型。
下面是分步解决的方案:
第一步:彻底修复col1的混合类型问题
首先必须统一col1的类型,这里推荐先把所有值转为字符串,再尝试安全转为浮点型(同时处理无法转换的异常值,避免后续报错):
from pyspark.sql import functions as F from pyspark.sql.types import FloatType # 先把col1强制转为字符串,统一格式 spark_df = spark_df.withColumn("col1_str", F.col("col1").cast("string")) # 尝试把字符串转为浮点型,无法转换的设为null(可根据业务需求调整默认值) spark_df = spark_df.withColumn( "col1", F.when( F.col("col1_str").rtrim(F.lit(".")).cast("int").isNotNull(), F.col("col1_str").cast(FloatType()) ).otherwise(F.lit(None).cast(FloatType())) ) # 清理临时列 spark_df = spark_df.drop("col1_str")
说明
这一步先统一所有col1的值为字符串格式,再通过判断是否能转换为数字来统一成浮点型。如果遇到非数字的字符串值,会被设为null,你可以根据业务逻辑替换成合适的默认值(比如0)。
第二步:正确执行重分区写入
解决类型问题后,再调整分区写入的逻辑:
- 不建议使用
repartition(1),这会把所有数据shuffle到一个分区,不仅性能极差,还容易因数据量过大导致OOM; - 确保分区列的类型统一(比如字符串类型,避免分区路径生成时出现异常);
- 如果需要控制输出文件数量,用
coalesce(n)替代repartition(n),它不需要shuffle,性能更优。
修改后的写入代码:
file_path_re = "s3://other_bucket/re-partition" partition_columns = ["source", "org_id", "device_id", "channel_id"] # 确保分区列都是字符串类型(可选,如果你确认它们类型已经统一) for col_name in partition_columns: spark_df = spark_df.withColumn(col_name, F.col(col_name).cast("string")) # 写入:先用overwrite测试,确认数据正确后再切换为append spark_df.coalesce(4) # 根据数据量调整文件数,比如4个 .write .partitionBy(partition_columns) .mode('overwrite') .parquet(file_path_re)
说明
- 使用
overwrite模式可以先清空目标路径的旧数据,方便验证数据是否正确写入; coalesce不会触发shuffle,适合在不改变分区分布的前提下调整输出文件数量;- 统一分区列类型可以避免生成的分区路径出现乱码或格式异常。
额外验证步骤
在写入前可以先确认数据修复结果,确保没有混合类型残留:
# 查看col1的类型分布 spark_df.select(F.col("col1"), F.typeOf(F.col("col1"))).distinct().show()
内容的提问来源于stack exchange,提问作者Vikram Ranabhatt
相关产品推荐
相关产品推荐

