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

Spark Repartition/Shuffle优化:预重分区耗时过长求优化方案

Spark Repartition操作耗时过长的优化方案

问题背景

在执行转换逻辑前对数据进行Repartition操作,但该过程耗时极长,结合集群配置、代码、提交参数及UI表现,给出以下优化建议。

集群配置

AWS EMR集群,200台任务节点r6g.16xlarge(64 vCore,488 GiB内存)

代码实现

df=spark.sql(''' 
        with basedata as (
            select kgs,marg_id,*******price
              from searcced_v2 
             where is_ot = 0 
             and pardate between "20240602" and "20240608"
             and rg in ('gg','MXg','gCA')  )
             select * from basedata ''').repartition(24000,'kgs','marg_id')
df.createOrReplaceTempView('vw_sewords')

Spark提交配置

spark-submit --conf spark.sql.files.maxPartitionBytes=268435456 \
--master yarn --deploy-mode cluster  --conf spark.yarn.maxAppAttempts=1 \
--conf spark.sql.adaptive.enabled=true --conf spark.dynamicAllocation.enabled=false \
--conf spark.sql.parquet.filterPushdown=true \
--conf spark.sql.adaptive.coalescePartitions.enabled=true \
--conf spark.sql.adaptive.advisoryPartitionSizeInBytes=268435456 \
--conf spark.sql.adaptive.optimizeSkewsInRebalancePartitions.enabled=true \
--conf spark.sql.adaptive.rebalancePartitionsSmallPartitionFactor=.5 \
--conf spark.sql.adaptive.coalescePartitions.parallelismFirst=false \
--conf spark.sql.adaptive.coalescePartitions.initialPartitionNum=24000 \
--conf spark.sql.adaptive.localShuffleReader.enabled=true \
--conf spark.shuffle.io.connectionTimeout=8000 \
--conf spark.network.timeout=50000s  --conf spark.files.fetchTimeout=600s \
--conf spark.serializer=org.apache.spark.serializer.KryoSerializer --conf spark.kryoserializer.buffer.max=1g \
--conf spark.memory.storageFraction=0.05 --conf spark.memory.fraction=.8 \
--conf spark.shuffle.compress=true --conf spark.shuffle.spill.compress=true \
--conf spark.hadoop.fs.s3.multipart.th.fraction.parts.completed=0.99 \
--conf spark.sql.objectHashAggregate.sortBased.fallbackThreshold=4000000 \
--conf spark.reducer.maxReqsInFlight=100 \
--conf spark.shuffle.io.retryWait=60s \
--conf spark.shuffle.io.maxRetries=10 \
--conf spark.reducer.maxSizeInFlight=1024m \
--conf spark.shuffle.file.buffer=1024k \
--conf spark.reducer.maxBlocksInFlightPerAddress=100 \
--conf spark.io.compression.codec=zstd \
--conf spark.shuffle.service.enabled=true \
--conf spark.io.compression.zstd.level=3 \
--conf spark.executor.cores=5 \
--conf spark.executor.instances=2400 \
--conf spark.executor.memory=34g --conf spark.driver.memory=60g --conf spark.executor.memoryOverhead=5g --conf spark.driver.memoryOverhead=4g \
--conf spark.hadoop.fs.s3a.fast.output.enabled=true \
--conf spark.executor.extraJavaOptions="-XX:+UseG1GC -XX:+UnlockDiagnosticVMOptions -XX:+G1SummarizeConcMark -XX:InitiatingHeapOccupancyPercent=35 -verbose:gc -XX:+PrintGCDetails -XX:+PrintGCDateStamps -XX:OnOutOfMemoryError='kill -9 %p' -Djavax.net.ssl.trustStore=/home/hadoop/.config/certs/InternalAndExternalTrustStore.jks" --conf spark.driver.extraJavaOptions="-XX:+UseG1GC " \
test.py

UI表现分析

截图显示该Repartition阶段存在明显数据倾斜:部分任务执行时间远超平均水平,Shuffle读写数据量巨大,这是导致耗时过长的核心原因;同时任务数设置较高,存在调度开销过大的可能。

优化方案

1. 根治数据倾斜

  • 先排查键值分布:执行SQL统计kgs和marg_id的数据量分布,找出异常大的键值对
    SELECT kgs, marg_id, COUNT(*) 
    FROM searcced_v2 
    WHERE is_ot = 0 
      AND pardate BETWEEN "20240602" AND "20240608" 
      AND rg IN ('gg','MXg','gCA')
    GROUP BY kgs, marg_id 
    ORDER BY COUNT(*) DESC 
    LIMIT 20
    
  • 针对倾斜键值处理:
    • 若为空值、无效值,直接过滤
    • 若为有效数据,采用加盐拆分法:给倾斜键添加随机后缀(如0-9),按加盐后的键分区处理,完成后再合并;或单独提取倾斜数据做特殊处理,最后与其他数据合并
  • 强化自适应倾斜优化:调整spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes(默认256MB),根据实际数据量设置合理阈值,让Spark自动识别并拆分倾斜分区

2. 合理调整分区数量

  • 当前24000个分区接近集群核数(200×64=12800)的2倍,处于合理范围,但如果存在大量小分区,可依赖spark.sql.adaptive.coalescePartitions.enabled自动合并
  • 取消手动指定分区数,改用repartitionByRange(若后续有排序需求),或直接移除repartition调用,让AQE自适应调整分区数量,减少手动配置的误差

3. 优化Shuffle性能

  • 调整Shuffle内存:若出现频繁磁盘Spill,可适当提高spark.shuffle.memoryFraction(默认0.2),配合spark.memory.fraction确保Shuffle有足够内存空间
  • 优化Shuffle IO:将spark.shuffle.unsafe.file.output.buffer从默认32KB提升至64KB/128KB,减少磁盘IO次数;保持Zstd压缩、Shuffle Service的启用状态,降低数据传输量

4. 优化数据源读取

  • 检查searcced_v2表的存储格式:确保Parquet/Orc的谓词下推生效,验证过滤条件是否真正下推到存储层
  • S3存储优化:设置spark.hadoop.fs.s3a.block.size与spark.sql.files.maxPartitionBytes匹配(256MB),减少小文件;启用spark.hadoop.fs.s3a.fast.upload=true提升读写速度
  • 提前分区:对searcced_v2表按pardate、rg做分区,减少每次读取的数据量

5. 集群资源调优

  • 调整Executor配置:当前每个节点运行12个Executor(2400/200),可将spark.executor.cores提升至8-10,减少Executor数量,降低调度开销;对应调整spark.executor.memory,保证每个Executor内存充足
  • 匹配YARN资源:检查YARN的节点内存、核数配置,确保与Spark参数匹配,避免资源浪费或不足

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 05:13:18