Spark中df.coalesce(1).write.parquet执行缓慢的原因及优化方法
PySpark写入Parquet到S3过慢问题排查与优化
问题场景
以下是一段PySpark代码,执行df.coalesce(1).write.parquet操作时耗时过长:
spark = SparkSession.builder \ .appName("MyApp") \ .master("local[*]") \ .config('spark.driver.host', '127.0.0.1') \ .getOrCreate() data = [('James','','Smith','1991-04-01','M',3000), ('Michael','Rose','','2000-05-19','M',4000), ('Robert','','Williams','1978-09-05','M',4000), ('Maria','Anne','Jones','1967-12-01','F',4000), ('Jen','Mary','Brown','1980-02-17','F',-1) ] columns = ["firstname","middlename","lastname","dob","gender","salary"] df = spark.createDataFrame(data=data, schema = columns) df.coalesce(1).write.parquet("s3a://sample/tesing.parquet",mode='overwrite')
导致耗时过长的原因
coalesce(1)强制单分区写入:该操作将所有数据合并到1个分区,所有写入任务只能由单个线程/Executor执行,完全丧失了Spark的分布式并行能力,数据量越大,串行处理的瓶颈越明显。- S3对象存储的特性限制:S3并非本地文件系统,单线程写入大文件时,无法利用多节点并行写入的优势,且S3单对象写入的吞吐量存在上限,单分区大文件写入会直接触发这个瓶颈。
- 本地模式的局限性:代码使用
master("local[*]")运行,本身并行度就远低于集群模式,再加上单分区的串行处理,进一步放大了速度问题。
优化方案
- 移除
coalesce(1),采用并行分区写入:
让Spark根据数据量自动管理分区,或者通过repartition(n)手动设置合适的分区数(建议每个分区大小控制在128MB-256MB,这是Parquet格式的最优分区尺寸),多个分区可以并行写入S3,大幅提升效率。 - 优化S3相关配置:
- 启用S3快速上传:添加配置
spark.hadoop.fs.s3a.fast.upload=true,开启多部分并行上传; - 调整分块大小:设置
spark.hadoop.fs.s3a.multipart.size=67108864(即64MB,可根据网络情况调整),减少分块上传的次数; - 增加S3连接数:设置
spark.hadoop.fs.s3a.connection.maximum=50(默认是15),提升并行连接能力。
- 启用S3快速上传:添加配置
- 切换到集群模式运行:如果处理的是大量数据,不要使用本地模式,将作业提交到Spark集群,利用多节点、多Executor的资源进行并行处理和写入。
- 检查预处理逻辑:确保写入前的数据处理步骤没有其他导致串行执行或数据倾斜的操作,保证整个作业的并行度。
内容的提问来源于stack exchange,提问作者Santhosh
相关产品推荐
相关产品推荐

