You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何优化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交互性能:
    --conf spark.hadoop.fs.s3a.committer.name=directory \
    --conf spark.hadoop.fs.s3a.fast.upload=true \
    --conf spark.hadoop.fs.s3a.connection.maximum=100
    
    这些配置能减少S3的重命名操作开销,启用多线程上传,增大连接数提升吞吐量。

资源配置优化

  • 调整并行度参数:将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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.05 11:15:14