如何优化PySpark DataFrame写入S3的上传性能
问题描述
使用PySpark通过JDBC从MySQL读取数据写入S3存储,实测单次拉取10万条记录耗时20-30秒,但同批次10万条数据写入S3需要5-7分钟。源库总数据量约3108700000条,按当前写入速度完成全量同步需要数月,已确认性能瓶颈仅存在于S3写入环节,读取环节无异常。
当前使用的代码如下:
df = spark.read.format("jdbc"). option('url', jdbcURL). option('driver', driver). option('user', user_name). option('password', password). option('query', data_query).load() output_df = df.persist() output_df.repartition(1).write.mode("overwrite").parquet(target_directory)
核心性能根因
当前写入慢的核心问题来自代码中的两处错误配置:
repartition(1)强制将全量数据通过Shuffle汇聚到单个Executor节点,完全放弃Spark分布式并行能力,单线程串行写入S3,受限于单节点网络带宽、S3客户端吞吐量上限,写入效率极低;同时会生成超大单文件,后续读取性能也会受影响。- 无复用场景下调用
persist(),会额外增加数据序列化、本地磁盘IO开销,无任何收益反而拖慢处理速度。
除此之外,未对S3客户端、Spark写入参数做针对性调优,也会在一定程度上影响写入效率。
可落地优化方案
- 移除
repartition(1)配置,合理设置写入并行度
按照Parquet单文件128MB-256MB的最佳实践设置分区数,结合集群总CPU核数调整,建议分区数设置为集群总Executor核数的2-3倍,让所有Executor节点并行写入S3,吞吐量可随集群资源线性提升。如果读取后的分区数过于零散,可使用coalesce调整到目标分区数,相比repartition可避免不必要的Shuffle开销。 - 删除无意义的
persist()操作
仅做单次写入的DataFrame不需要持久化,删除该行代码即可减少额外的IO开销。 - 调优S3客户端与写入参数
初始化SparkSession时添加如下配置,提升S3客户端上传效率:spark = SparkSession.builder \ .config("spark.hadoop.fs.s3a.threads.max", "200") # S3客户端并发上传线程数,按集群规模调整 .config("spark.hadoop.fs.s3a.connection.maximum", "500") # S3客户端最大HTTP连接数 .config("spark.hadoop.fs.s3a.fast.upload", "true") # 开启快速上传,无需等数据全量缓存到本地再上传 .config("spark.hadoop.fs.s3a.multipart.size", "128M") # S3分片上传单分片大小 .config("spark.hadoop.fs.s3a.multipart.threshold", "128M") # 触发分片上传的文件大小阈值 .config("spark.hadoop.mapreduce.fileoutputcommitter.algorithm.version", "2") # 启用V2版本提交算法,大幅减少S3上的重命名操作开销 .config("spark.sql.parquet.compression.codec", "snappy") # 使用Snappy压缩,压缩解压速度快,减少网络传输数据量 .getOrCreate() - JDBC读取阶段直接设置并行度,避免后续Shuffle
读取MySQL时直接添加分区参数,让数据读取阶段就均匀分布在所有Executor节点,写入时不需要做额外的Shuffle重分区:- 选择数值型、分布均匀的列作为分区列(比如自增主键、创建时间戳)
- 配置
lowerBound、upperBound指定分区列的取值范围 - 配置
numPartitions为目标并行度,和写入需要的分区数保持一致
- 适当提升单批次处理数据量
当前单批次10万条的量级偏小,任务调度、初始化的固定开销占比过高。建议单批次数据量按每个分区对应128MB-256MB数据调整,单批次总数据量可提升到千万级,减少总批次数量,降低调度开销。
优化后参考代码
# 初始化SparkSession加载S3优化配置 spark = SparkSession.builder \ .appName("MySQL-to-S3-Sync") \ .config("spark.hadoop.fs.s3a.threads.max", "200") \ .config("spark.hadoop.fs.s3a.connection.maximum", "500") \ .config("spark.hadoop.fs.s3a.fast.upload", "true") \ .config("spark.hadoop.fs.s3a.multipart.size", "128M") \ .config("spark.hadoop.mapreduce.fileoutputcommitter.algorithm.version", "2") \ .config("spark.sql.parquet.compression.codec", "snappy") \ .getOrCreate() # JDBC读取时直接设置并行度,以自增主键id为分区列举例 df = spark.read.format("jdbc") \ .option('url', jdbcURL) \ .option('driver', driver) \ .option('user', user_name) \ .option('password', password) \ .option('query', data_query) \ .option('partitionColumn', 'id') \ .option('lowerBound', 1) \ .option('upperBound', 3108700000) \ .option('numPartitions', 200) # 按集群总核数调整 .load() # 直接写入,无额外缓存、重分区操作 df.write.mode("overwrite").parquet(target_directory)
按上述方案优化后,若集群资源足够,写入速度可提升数十到上百倍,31亿条数据的同步周期可从数月压缩到数小时到数天级别。
内容的提问来源于stack exchange,提问作者Yogesh Sharma
相关产品推荐
相关产品推荐

