如何优化PySpark批量写入S3 Parquet文件的处理速度?
PySpark写入S3性能优化需求
我需要用PySpark处理每个含70万行的Parquet文件并写入S3,当前单次流程(读取到写入)平均耗时80秒,其中仅写入就占50秒。后续要处理约6000个DataFrame,还要循环执行1万次,希望把单次流程耗时降到30秒左右。
现有代码片段
original_dataframe = spark.read.option('header', 'true').option('inferSchema','true').parquet(f"s3://{bucketname}/{data_files_path}") distinct_df = original_df.distinct() distinct_df_cache = distinct_df.cache() distinct_df_cache.coalesce(8).write.format("parquet") .option("header","true").option("inferSchema","true") .mode("overwrite").save(f"s3://{bucketname}/{data_files_path}")
现有Spark提交命令
sudo spark-submit --master yarn --executor-memory 48G --driver-memory 30g --num-executors 8 --executor-cores 8 --conf spark.sql.shuffle.partitions=16 --conf spark.sql.parquet.mergeSchema=true --conf spark.sql.parquet.filterPushdown=true --conf spark.default.parallelism=90 --conf spark.sql.parquet.writeLegacyFormat=true my_Code.py
优化方案
读取阶段优化
- 移除
inferSchema参数:Parquet文件自带Schema,启用inferSchema会额外扫描文件推断结构,直接使用spark.read.parquet()即可,能大幅减少读取耗时。 - 列裁剪与分区过滤:仅读取业务需要的列(用
select()指定),如果数据按分区存储,通过where()过滤目标分区,减少读取的数据量。
数据处理阶段优化
- 替换
distinct()为dropDuplicates():两者功能一致,但dropDuplicates()可指定具体去重列,减少计算开销;若必须全列去重,调整shuffle分区数匹配资源。 - 删除不必要的
cache():当前流程是读取→去重→写入,没有重复使用缓存的DataFrame,cache()反而会增加内存/磁盘IO开销,直接移除即可。 - 调整shuffle分区数:将
spark.sql.shuffle.partitions从16改为64(等于executor数×executor cores),让shuffle数据分布更均衡,避免单个分区过大拖慢处理速度。
写入阶段优化
- 用
repartition(8)替代coalesce(8):coalesce是合并分区,容易导致数据倾斜;repartition会重新均匀分配数据,写入S3时每个文件大小更合理,提升写入效率。 - 移除写入时的
header和inferSchema参数:Parquet格式不需要header,写入时指定inferSchema完全多余,删除这两个参数减少不必要的处理。 - 添加S3读写优化配置:在Spark提交命令中增加以下配置,提升S3交互性能:
这些配置能减少S3的重命名操作开销,启用多线程上传,增大连接数提升吞吐量。--conf spark.hadoop.fs.s3a.committer.name=directory \ --conf spark.hadoop.fs.s3a.fast.upload=true \ --conf spark.hadoop.fs.s3a.connection.maximum=100
资源配置优化
- 调整并行度参数:将
spark.default.parallelism从90改为128(executor数×executor cores×2),让任务并行度与集群资源匹配。 - 降低driver内存:将
--driver-memory从30G调整为16G,若没有大量元数据处理需求,多余内存可留给executor提升计算能力。 - 启用动态资源分配:添加
--conf spark.dynamicAllocation.enabled=true,让Spark根据任务负载自动调整executor数量,避免资源闲置或不足。
内容的提问来源于stack exchange,提问作者kislay kashyap
相关产品推荐
相关产品推荐

