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
相关产品推荐
相关产品推荐

