PySpark合并Parquet文件写入S3报解码错误解决方案
问题背景
已成功读取S3路径s3://aws-emr-resources-359367213591-us-east-1/taxi_data_2020/下的12个小型Parquet文件并完成合并,经统计总数据量为24649092条。通过coalesce(1)将合并后的DataFrame合并为单个分区,尝试以单个Parquet文件形式写入目标S3路径s3://aws-emr-resources-359367213591-us-east-1/merged_data_2020/single时任务失败抛出异常。
复现代码
SparkSession初始化
from pyspark.sql import SparkSession spark = SparkSession.builder.appName('combine_files').getOrCreate()
源文件读取
df=spark.read.parquet("s3://aws-emr-resources-359367213591-us-east-1/taxi_data_2020/*").coalesce(1) df.count() # 执行结果:24649092
写入操作
df.write.parquet("s3://aws-emr-resources-359367213591-us-east-1/merged_data_2020/single")
核心报错信息
任务抛出
org.apache.spark.SparkException: Job aborted异常,根因为java.lang.UnsupportedOperationException: org.apache.parquet.column.values.dictionary.PlainValuesDictionary$PlainDoubleDictionary,Parquet读取流程中尝试调用decodeToInt方法解码Double类型字典值,任务重试4次后最终失败。
问题原因
- 核心触发点为源Parquet文件Schema不一致:12个待合并文件存在同名字段类型定义冲突,部分文件中同名字段为Double类型,其余文件中该字段为整数类型。Spark自动推断Schema读取时虽做了表层类型兼容,但
coalesce(1)触发全量数据拉取到单节点、写入阶段解码Parquet字典页时,类型映射逻辑出错,对Double类型的字典值调用了仅支持整数类型的decodeToInt方法,触发不支持操作异常。 - 额外稳定性隐患:2400余万条数据直接调用
coalesce(1)会将全量数据的计算、序列化、写入压力全部集中到单个Executor节点,即使绕过当前类型报错,也极大概率触发内存溢出、GC超时、磁盘溢写失败等问题导致任务失败。
解决方法
- 显式指定读取Schema,从根源避免自动推断带来的类型冲突:读取前根据业务规则明确定义所有字段的类型,将存在类型冲突的数值字段统一声明为
DoubleType,避免Spark自动类型推断的逻辑漏洞。示例代码如下:
from pyspark.sql.types import * # 按实际表结构补全所有字段定义,存在类型冲突的数值字段统一设为DoubleType custom_schema = StructType([ # 示例:StructField("trip_distance", DoubleType(), True), StructField("total_amount", DoubleType(), True) ]) df = spark.read.schema(custom_schema).parquet("s3://aws-emr-resources-359367213591-us-east-1/taxi_data_2020/*")
- 调整分区逻辑,规避单节点计算压力:不要在读取后直接调用
coalesce(1),先通过repartition将数据打散到合理数量的分区完成类型校验、数据处理,写入阶段再合并为单分区,降低单点负载。示例代码如下:
# 按集群资源情况设置分区数,通常按每个Executor分配2-3个分区设置 df = df.repartition(48) # 写入前再合并为单分区输出单个Parquet文件 df.coalesce(1).write.parquet("s3://aws-emr-resources-359367213591-us-east-1/merged_data_2020/single")
- 兜底兼容方案:若上述调整后仍存在字典解码报错,可在初始化SparkSession时关闭Parquet向量化读取配置,绕开字典解码的类型冲突逻辑:
spark = SparkSession.builder.appName('combine_files') \ .config("spark.sql.parquet.enableVectorizedReader", "false") \ .getOrCreate()
内容的提问来源于stack exchange,提问作者Emislim
相关产品推荐
相关产品推荐

