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

为何增加Spark作业核心数未缩短耗时?分区已均匀分布

PySpark Parquet 读写性能未达预期问题分析

用户代码

spark = SparkSession.builder \
    .appName("Read Write Parquet EFS") \
    .getOrCreate()
input_folder_path = "/home/ubuntu/efs-spark/10M-32part/input"
output_folder_path = "/home/ubuntu/efs-spark/10M-32part/output"
df = spark.read.parquet(input_folder_path)
df.write.parquet(output_folder_path)

测试情况

  • 16核CPU:总耗时10秒,输入分区数16(单分区625000行),写入耗时5秒,单任务耗时4000-5000ms
  • 32核CPU:总耗时10秒,输入分区数32(单分区312500行),写入耗时5秒,单任务耗时4000-5000ms

核心原因与排查方向

  1. EFS存储系统瓶颈
    EFS的吞吐量存在上限,分为突发模式和预置吞吐量模式。若使用突发模式,耗尽突发额度后吞吐量会降至基准值,此时16核场景已达EFS吞吐量饱和点,32核下新增的任务仅会加剧IO竞争,单个任务的等待时间抵消了数据量减半的优势,最终总耗时不变。

  2. Spark资源配置未充分利用32核
    默认Spark配置不会自动适配全部CPU核心,需检查:

    • 若为local模式,需指定local[32]而非默认local[*],避免JVM限制未完全利用核心
    • 集群模式下,确认spark.executor.cores、spark.executor.instances、spark.total.executor.cores是否配置为匹配32核的资源,确保任务能真正并行执行
  3. EFS并发连接限制
    EFS对并发写入连接数有隐性限制,过多并行写入任务会导致连接排队,每个任务的IO等待时间延长,使得单任务耗时未随数据量减少而降低。可尝试减少分区数(如24)测试总耗时变化,或升级EFS至预置吞吐量模式提升总带宽。

  4. Parquet写入配置优化
    确认spark.sql.parquet.enableVectorizedWriter已开启(默认开启),该配置能提升Parquet写入的CPU效率;同时检查压缩 codec(如snappy)是否最优,避免序列化/压缩占用过多CPU资源,但此点并非本次性能瓶颈的核心原因。

内容的提问来源于stack exchange,提问作者kishore

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 18:43:28