如何减少PySpark写入Parquet文件的耗时?
PySpark写入Parquet文件耗时优化方案
你在PySpark中完成7143779条记录的转换仅耗时45秒,但写入Parquet却花费5分35秒,可通过以下方法优化写入流程:
1. 调整数据分区,优化并行写入效率
默认的分区数可能不匹配数据量,导致写入时任务不均衡或产生过多小文件。可以根据集群资源调整分区数,让每个分区数据量控制在30-50万条左右:
# 在写入前调整分区 final_df = final_df.repartition(20) # 20为示例值,可根据实际情况调整
也可以通过Spark配置全局调整shuffle分区数:
spark.conf.set("spark.sql.shuffle.partitions", "20")
如果数据分区过大但无shuffle需求,也可使用coalesce合并分区(注意数据均衡性)。
2. 启用高效压缩算法,减少IO开销
Parquet支持多种压缩算法,启用压缩能大幅降低磁盘IO量,提升写入速度。推荐使用Snappy(速度与压缩比平衡)或ZSTD(更高压缩比):
df.write.mode('overwrite').option("compression", "snappy").parquet(path)
3. 移除重复的count()操作,避免额外计算
当前写入函数中调用了两次df.count(),这会触发全量数据扫描的Action操作,完全是冗余计算。可以提前计算记录数并传入函数:
修改调用代码:
record_count = final_df.count() func.write_df_as_parquet_file(final_df, cur_path, logger, record_count)
修改写入函数:
def write_df_as_parquet_file(df, path, logger, record_count): try: df.write.mode('overwrite').parquet(path) logger.debug(f'File written Successfully at {path} , No of records in the file : {str(record_count)}') print(f'File written Successfully at {path} , No of records in the file : {str(record_count)}') except Exception as exc: return_code = 'file Writting Exception: ' + path + '\n' + 'Exception : ' + str(exc) print(return_code) logger.error(return_code) raise
4. 按业务字段分区写入(若适用)
如果数据存在合适的分区键(如时间、类别等),可通过partitionBy按字段分区写入,既能提升后续查询效率,也能让写入任务并行处理不同分区的数据:
df.write.mode('overwrite').partitionBy("partition_column").parquet(path)
注意:分区键的基数不宜过大,否则会生成大量小文件,反而降低性能。
5. 优化Spark IO相关配置
调整以下Spark配置参数,减少写入时的元数据开销与资源瓶颈:
- 关闭不必要的元数据生成:
spark.conf.set("spark.hadoop.parquet.enable.summary-metadata", "false") - 禁用旧格式兼容,使用新Parquet格式:
spark.conf.set("spark.sql.parquet.writeLegacyFormat", "false") - 根据集群资源调整executor资源:增加
spark.executor.cores和spark.executor.memory,提升并行处理能力,减少GC停顿。
6. 优化Union后的分区均衡性
如果df1和df2的分区数差异较大,Union后的数据分区可能不均衡,导致写入时部分任务耗时过长。可以在Union后执行repartition调整分区,保证每个分区的数据量均匀。
内容的提问来源于stack exchange,提问作者user3521180
相关产品推荐
相关产品推荐

