Spark新手咨询:模糊匹配后PySpark写入Parquet速度极慢问题
问题分析与优化建议
结论:当前I/O耗时不正常
30MB级别的数据写入Parquet,无论本地还是GCP集群,3分钟的耗时都远超合理范围,核心原因是配置不合理导致I/O开销被过度放大。
核心问题与优化方案
1. 写入批次设置过小
你设置的spark.sql.parquet.write.batchSize=100太小,这会让Spark频繁执行小批量写入操作,每个批次的文件打开、关闭、元数据同步等I/O开销占比极高,直接拖慢整体速度。
- 优化操作:将batchSize调整为默认值
10000或更高,示例代码:sample_spark.write.format("parquet") \ .mode("overwrite") \ .option("compression", "lz4") \ .option("spark.sql.parquet.write.batchSize", "10000") \ .save("parquet_5k")
2. 并行度与分区数不匹配
本地模式:你设置了
spark.default.parallelism=300,但本地只有7核,300个并行任务会导致严重的任务碎片化——每个任务处理的数据量极少,大部分时间浪费在任务调度和小文件I/O上。- 优化操作:将并行度调整为核心数的1-2倍(比如
14),同时写入前用coalesce(7)合并分区,减少小文件数量:spark = SparkSession.builder \ .appName("Spark NLP") \ .master("local[7]") \ .config("spark.driver.memory", "16g") \ .config("spark.default.parallelism", "14") \ .config("spark.jars.packages", "com.johnsnowlabs.nlp:spark-nlp_2.12:3.1.0") \ .config("spark.kryoserializer.buffer.max", "1000M") \ .getOrCreate().newSession() # 合并分区后写入 sample_spark.coalesce(7).write.format("parquet")... - 额外说明:本地模式下
spark.executor.memory、spark.executor.cores、spark.executor.instances、spark.dynamicAllocation.enabled这些配置均不生效(本地仅存在Driver进程),可以直接删除以避免混淆。
- 优化操作:将并行度调整为核心数的1-2倍(比如
GCP集群模式:集群总核数为
7*10=70,spark.default.parallelism=300同样过高,建议调整为70-140之间,同时用repartition(70)设置合理的分区数,让每个executor核心对应一个分区的写入任务,避免资源浪费。
3. 不必要的读取配置干扰
你设置了spark.sql.parquet.enableVectorizedReader=false,这个参数仅用于优化Parquet读取性能,与写入无关,且会影响后续读取该Parquet文件的效率。写入时无需设置此参数,保持默认值(true)即可。
4. 存储介质与网络因素
- 本地运行:如果使用机械硬盘(HDD),I/O速度远慢于SSD,建议更换为SSD存储以提升写入速度。
- GCP集群:检查存储区域是否与集群区域一致,跨区域存储会增加网络I/O延迟;同时可考虑使用高性能存储类型(如SSD持久化磁盘)。
优化后预期效果
调整上述配置后,5000行样本的Parquet写入时间应能压缩至几十秒以内,样本量增大后耗时也会保持线性增长,而非指数级增加。
内容的提问来源于stack exchange,提问作者Sam Ding
相关产品推荐
相关产品推荐

