分布式计算框架优化实战:从数据倾斜到小文件治理

发布时间:2026/10/9 12:40:26
分布式计算框架优化实战:从数据倾斜到小文件治理 上个月我们团队接到一个棘手的任务核心报表链路每天跑一次耗时从年初的半小时慢慢涨到了两个多小时业务方每天早上都在群里催数。一开始大家习惯性想加资源但加了executor、扩了队列效果也就那样成本倒是实打实上去了。后来我从头到尾把这条链路翻了一遍从数据分布、执行计划、框架参数到存储格式逐一排查最后把整条作业链压缩到25分钟左右跑完还能把资源让回给别的团队。这篇文章就把这套分布式计算框架优化的思路和操作细节整理出来希望能帮到正在被慢作业折磨的工程师。这篇文章不聊虚的重点讲清楚几件事怎么快速定位瓶颈、数据倾斜怎么治、资源参数怎么配才合理、SQL怎么写才不浪费计算资源、小文件怎么合并以及一整条链路优化的先后顺序。适合用Spark、Hive、Flink做离线或实时计算的同学参考尤其是那种作业能跑但越跑越慢的团队。1. 先定位瓶颈优化动作前必须搞清的问题很多人一上来就调参数、改SQL结果改了一通不知道到底改对了没有。我自己的经验是先花半小时把执行链路看清楚尤其把瓶颈定位准确后面改起来才有方向改完才知道是不是真的生效。1.1 Spark UI里最该看的三块信息无论你用Spark、Flink还是跑在Yarn上的MapReduce第一件事都是打开作业的UI页面。以Spark为例重点看三块Stage列表每个Stage的耗时和Shuffle读写量。如果Stages很多、依赖链很长先用DAG图看整体结构找出哪几个Stage占据了绝大部分耗时。Task耗时分布点进某个耗时的Stage后观察所有Task的耗时。如果大量Task在秒级完成极少数Task要跑几十分钟基本可以断定是数据倾斜如果所有Task耗时都比较均匀地长那可能是数据量大、资源不足或者计算逻辑本身太重。Executor的GC时间与内存曲线如果Event Timeline上能看到频繁的Full GC并且GC时间占比超过10%说明内存堆配置可能有问题或者某个算子把大量数据堆积在内存里。举个例子我接手的那条链路里有一个Join StageUI上显示6000个Task其中5900个在1到3秒内完成剩下100个都在10分钟以上。这时候看Executor的CPU利用率整体其实不高问题明显不是资源不够而是少数Task拖慢了整个Stage。1.2 不同瓶颈的判读标准定位瓶颈时最怕的就是把数据倾斜和资源不足搞混。下面的表格是我平时做快速判读用的分享出来供参考现象优先怀疑方向验证方法少量Task耗时远超中位数且这些Task处理的数据量如Shuffle Read Bytes明显偏大数据倾斜检查Stage输入数据的Key分布所有Task耗时同步变长且CPU/内存利用率高资源不足或数据量整体增长看集群当前是否有其他作业抢占资源Task耗时不高但Stage总数很多依赖链长执行计划不合理用explain看逻辑计划和物理计划GC时间占比高任务反复出现内存溢出内存参数或算子内存占用问题调整Executor内存与spark.memory参数磁盘IO或网络IO接近饱和基础设施瓶颈或Shuffle量过大查看节点监控确认数据本地性这一步骤的核心是用数据说话。不要凭感觉觉得某个SQL慢就是Join慢先看Stage和Task分布再看Shuffle量顺藤摸瓜找到真正的问题点。2. 数据倾斜治理分布式计算里性价比最高的优化点数据倾斜是分布式计算框架优化中最常见也最典型的性能杀手。我见过的生产事故里十个慢作业里有六七个都跟数据倾斜有关。2.1 倾斜是怎么发生的分布式计算的核心思想是把数据分到多个Task并行处理。Shuffle阶段会按照Key的哈希值对数据进行分区落到不同Task上。问题在于如果某个Key在数据集中占据的份额过大比如某个热门商品ID的订单量是其他商品的上百倍那个处理该Key的Task就要比别人多算几十倍的数据于是整个Stage的耗时被它所拖累。常见的倾斜Key一般有这么几类空字符串、默认值比如“default”、“unknown”、特定业务ID如测试账号、超大店铺这些Key在业务上往往没什么聚合价值却在计算时造成了严重倾斜。2.2 三个立竿见影的治理方案方案一过滤或剔除无效Key如果倾斜来自空值或默认值且这些数据本身对结果无意义最简单的方式就是提前过滤。比如统计订单量时把user_id为空或为默认值的数据直接去掉SELECT user_id, COUNT(*) AS cnt FROM orders WHERE user_id IS NOT NULL AND user_id ! GROUP BY user_id;方案二大表Join小表时用广播业务里最常见的是大表关联维度表。维度表往往只有几千到几万行但是无论表多小如果不做任何优化Spark默认都会走Shuffle。正确的做法是让Executor把小表加载到内存里直接在每个Task内部完成关联彻底避免Shuffle。用Spark SQL可以这样SELECT /* BROADCAST(dim) */ a.id, b.name, a.amount FROM fact_table a JOIN dim_table b ON a.dim_id b.id;也可以动态判断表的大小Spark 3.x中如果一张表小于spark.sql.autoBroadcastJoinThreshold默认10MB会自动走BroadcastJoin不需要额外加Hint。但要注意当两张参与Join的表都很大时这条路走不通。方案三加盐拆分两阶段聚合如果倾斜的Key既有业务价值又无法过滤最常用的手段是加盐Salting。原理是把一个倾斜的Key加上随机前缀拆成多个子Key让数据分散到不同Task上先做一轮局部聚合再去掉前缀做全局聚合。# 伪代码两阶段聚合 from pyspark.sql import functions as F, Window # 第一轮倾斜key加随机后缀打散 df_expanded skewed_df.withColumn( salt_key, F.concat(F.col(key), F.lit(_), F.rand() * 10) ) partial_agg df_expanded.groupBy(salt_key).agg(F.sum(amount).alias(partial_amount)) # 第二轮去盐后做全局聚合 final_result ( partial_agg.withColumn(raw_key, F.split(salt_key, _)[0]) .groupBy(raw_key) .agg(F.sum(partial_amount).alias(total_amount)) )加盐粒度决定了效果盐分越多倾斜Key被打散得越均匀但第二轮聚合时需要重读一轮数据。实际调优时我一般先尝试10到100之间的盐值观察Task耗时分布再逐步调整。2.3 治理效果怎么验证改完之后不要只看作业总耗时还是回到Stage的Task耗时分布图去看。如果之前有Task跑30分钟现在所有Task都在3分钟内完成说明倾斜确实被解决了。另外要留意Shuffle Read Bytes的变化确认是不是真的把某一批超大Key的数据拆开了而不是把问题从某个Task转移到了另一个Task。3. 资源参数与调度策略从“加机器”到“用对参数”在没定位瓶颈之前盲目加资源是分布式计算框架优化中最容易踩的坑。但反过来当确认了作业确实存在资源不足时怎么配置参数又成了另一门学问。3.1 Executor数量、内核与内存的配比逻辑Executor的数量和规格决定了作业有多少并行度、多少内存可用。这里容易犯的毛病是“越多越好”。实际上Executor内存过大或过小都会出问题内存过大GC时停顿时间也长内存过小数据在内存里放不下频繁落盘。我常用的配比参考如下假设单节点可用资源比较充裕场景每个Executor核数每个Executor内存说明中小型ETL作业24GB避免浪费资源任务调度更细数据量大但逻辑简单48GB平衡吞吐与GCJoin/聚合密集412GB以上预留Shuffle读缓存同时要关注Spark内存模型。Executor内存被分为执行内存和存储内存通过spark.memory.fraction默认0.6控制。如果作业Shuffle量很大可以适当调高执行内存的比例如果Cache的表很多则维持默认或调低。具体数值要结合GC表现来验证不要照抄网上的配置。3.2 Shuffle并行度与数据量的匹配Shuffle并行度是个经常被忽视但影响巨大的参数。Spark里有两个关键参数spark.sql.shuffle.partitionsSpark SQL执行Shuffle时默认分区数默认值是200spark.default.parallelismRDD操作默认的并行度如果一份数据量很大但并行度没提上去相当于每个Task处理的数据量过大反过来并行度太高又会增加调度和元数据开销。比较通用的经验是让每个分区处理的数据量在128MB到256MB之间。举例来说如果Shuffle读的总数据量是200GB那么把spark.sql.shuffle.partitions设为1000到1500之间比较合理每个分区大概处理130到200MB。-- 通过spark-submit参数设置 spark-submit --conf spark.sql.shuffle.partitions1200 ...3.3 队列与动态资源分配的实际取舍在生产集群里多团队共用一个资源池很常见。如果一个大作业把队列打满其他团队的作业全被卡死这种“优化”得不偿失。我建议开启动态资源分配让作业根据实际任务量伸缩Executorspark.dynamicAllocation.enabledtruespark.dynamicAllocation.minExecutors2spark.dynamicAllocation.maxExecutors50spark.dynamicAllocation.executorIdleTimeout60s开启后作业两头都有余量数据量大时自动扩展数据量小的时候自动回收。但要注意动态分配和某些离线调度系统配合时可能因为初始化过慢而导致延迟需要在调度器层面做兼容。如果你们平台对资源有严格配额也可以关闭动态分配改用固定值但一定要在业务低峰期运行大作业。4. SQL代码层面的执行计划优化参数调整和数据倾斜解决了大部分问题但还有一种情况框架运行“健康”数据分布也均匀可作业还是很慢。这时候要回到SQL本身看执行计划。4.1 用Explain看执行计划差异无论Spark SQL还是Hive跑之前都可以先看执行计划。Spark 3.x里执行计划分为逻辑计划和物理计划重点看物理计划中的Join算法BroadcastHashJoin小表广播无Shuffle最快SortMergeJoin两侧都需要Shuffle和排序数据量大时必然出现EXPLAIN SELECT /* BROADCAST(dim) */ a.id, b.name FROM fact a JOIN dim b ON a.dim_id b.id;我在实际优化中看过不少案例大表和小表Join因为小表稍微超过了广播阈值结果走了SortMergeJoinShuffle量从2GB变成30GB耗时翻了好几倍。遇到这类情况要么给SQL加BroadcastHint要么调大spark.sql.autoBroadcastJoinThreshold但阈值不要设太大否则小表广播时GC压力会增大。4.2 谓词下推与列裁剪执行计划优化的另一重点是“过滤器下推”。SQL里写WHERE条件时框架不一定能保证第一时间过滤数据。下面是两个常见的处理原则先过滤再Join把数据量大的表先做WHERE过滤再做关联让Shuffle的数据量尽量小只取需要的列避免SELECT *减少列裁剪的数据量虽然Spark SQL的优化器会自动做谓词下推但有些场景下还是需要手动创造条件。比如子查询里关联条件写的时机不对或使用了UDF都可能阻碍下推。4.3 常见SQL写法坑整理几个自己在生产环境遇到过高频问题Avoid GroupByKey类操作在RDD API里groupByKey会直接把同一个Key的所有Value集中到一个Task上再处理。哪怕数据不倾斜GroupByKey的Shuffle量通常也比ReduceByKey大一截。能先聚合再传上游的尽量先做一轮合并。# 不好groupByKey后逐条求和 rdd.groupByKey().mapValues(lambda vals: sum(vals)) # 好reduceByKey直接在shuffle前合并 rdd.reduceByKey(lambda a, b: a b)Count Distinct的陷阱COUNT(DISTINCT col)在数据量巨大时往往需要精确去重Shuffle和内存开销非常高。如果业务上能接受近似值用approx_count_distinct会走HyperLogLog算法性能和准确性之间的平衡很不错。不等值连接Join条件中如果包含a.x b.x这类不等值条件无法走Broadcast Join只能做笛卡尔积慢是必然的。尽量改造成等值连接或者分区间处理再合并。执行计划优化没有银弹关键是每改一个点就跑一遍Explain看差异把SQL的执行路径和Shuffle数据量对比一下才会越来越顺手。5. 小文件治理与存储格式选择被低估的稳定性问题分布式计算框架优化的最后一环往往不在计算本身而在于计算之后写出来的数据。小文件问题初期不那么明显但会慢慢拖垮整个链路。5.1 小文件为什么可怕一个100GB的数据集如果以1MB的小文件存储会有10万个文件。每个文件在读取时都要经过NameNodeHDFS或对象存储的元数据查询在高并发下元数据服务会成为瓶颈同时任务在调度时也会因为分片过多产生大量空转Task。最直接的表现就是作业启动时间变长了处理同样的数据效率却逐年下降。造成小文件的常见场景包括动态分区插入时每个分区写入的并发Task数过多流式任务频繁写入同一个表Hive表经过多次INSERT OVERWRITE但Spark/Hive没做合并5.2 几种合并小文件的实操方案方案一写数据前控制分区数如果业务上对写入的分区数量没有严格要求可以在写入前先做一次重分区// 例如适合32个左右分区 df.repartition(32).write.mode(overwrite).partitionBy(dt).parquet(outputPath)方案二定时合并任务对于经常动态写入的表写一个定时的合并作业读取数据后重新写入。spark.sql(SET spark.sql.shuffle.partitions200) val df spark.read.parquet(inputPath) df.coalesce(50).write.mode(overwrite).parquet(outputPath)这里coalesce适合在数据量不大时减少分区但如果数据量很大用repartition配合sortBy来保证合并后文件大小均匀。方案三借助小文件合并工具现在很多数据平台都有专门的小文件合并工具原理可以简单理解成扫描小文件按目标大小合并成更大的文件然后更新表的元数据。有工具直接用工具没有的话参照上面的方案写个定期任务也能解决问题。5.3 存储格式与压缩的选型建议小文件合并完还要选对存储格式。对于分布式计算框架来说我强烈建议列式存储。列式存储Parquet、ORC在查询时只需读取需要的列对分析型作业收益非常大同时这类格式自带统计信息可以做谓词下推、更高效的压缩读同一份数据的IO开销能小很多。压缩格式方面生产环境我用得比较多的组合是存储格式压缩适用场景ParquetSnappy通用分析、读取较多ORCZSTDHive作业多、追求更高压缩率ParquetZSTD数据量非常大且读取频繁ZSTD的压缩率和解压速度都不错但要注意不同版本的Spark对它的支持程度老版本可能需要额外引入Jar包。选压缩格式时不能一味追求“压缩率最高”还要考虑压缩和解压的CPU开销。实际测试下来Snappy在绝大多数场景下的综合表现最稳。6. 优化链路实测一套组合拳把作业从2小时压到25分钟最后用一个完整的案例把前面的方法论串起来。这条链路当时是我们最核心的报表前置任务主要做的是大表关联维表、按维度聚合、再写回结果表。6.1 优化前的情况优化前作业总耗时123分钟整体状况如下检查项实际表现Stage耗时Top1一个Join Stage占87分钟Task耗时分布100个Task超过30分钟其余都在秒级完成Shuffle总量36GB小文件数量结果表写入了5000多个小文件队列情况固定Executor 100个资源利用率很低6.2 一步步做了什么第一步看执行计划发现一个事实表和维度表的Join走了SortMergeJoin而且Seen表过滤条件写得比较靠后导致Shuffle数据量过大。第二步调整广播Join维度表只有两万行左右加上BROADCASTHint后Join Stage的Shuffle量从36GB降到4GB。这一步直接把两个多小时的作业压缩到70分钟。第三步处理倾斜Key聚合阶段发现部分维度出现倾斜采用加盐两阶段聚合将100个长尾Task控制在2分钟以内。这一步之后作业总耗时降到38分钟。第四步调整并行度和Executor规格把不合理的Executor数量从100降到40每个Executor从4核8GB调整为6核12GB同时将spark.sql.shuffle.partitions从默认的200改为400GC时间明显下降。第五步小文件合并与写回优化重跑作业前检查了目标表目录发现大量小文件。先做了历史数据合并并在写入时采用repartition(64)控制写文件数量。6.3 整个过程的耗时变化阶段作业总耗时说明优化前123分钟排序连接倾斜小文件广播Join后70分钟消除大量Shuffle加盐聚合后38分钟消除了Task长尾参数调整后25分钟并行度和内存更匹配小文件合并25分钟稳定写结果表更快、后续读取也更快6.4 后续维护的经验优化不是一次性的。业务数据分布会变化新功能上线后SQL逻辑也常调整最关键的是要建立“基线告警”机制给核心作业设定基线耗时超过120%或者出现新增长尾Task就触发告警提醒。另外我倾向于每季度对核心表的数据量、文件数、分区数做一次体检把问题消灭在出现“明显变慢”之前。如果团队条件允许还可以把关键作业的Spark UI指标自动采集下来比如每阶段的Shuffle量、Task耗时中位数、GC时间存成历史报表。这样以后作业再变慢可以直接对比历史数据快速判断是新逻辑引入的问题还是数据量自然增长导致的。最后再分享一点自己的体会做分布式计算框架优化这几个月我最大的感受是顺序很重要。先看数据分布再看执行计划然后才是参数调整和存储治理。如果一上来就动参数很容易把问题掩盖在一堆变了的配置里后面再排查反而更难。还有一个容易被忽略的点是优化之前一定要记录基准数据改完每一步都回看指标否则你会分不清到底是因为哪一步变快了。希望这篇文章里的定位方法和组合拳能帮你在自己的集群上少走一些弯路。

关于本文作者

来自尧图内容编辑团队

尧图内容编辑团队 内容团队

尧图内容编辑团队

本文由尧图网络内容编辑团队执笔。团队由资深项目经理、前端工程师与设计师组成,所有内容均来自亲手交付的真实项目,先讲清问题、再给出可落地的解法。尧图深耕北京网站建设十年,服务过京华建材集团、智造科技等各行业客户,把一线经验沉淀为可复用的行业观察。

  • 十年建站经验,覆盖建材、制造、服务、文创等
  • 项目经理把关选题与事实准确性
  • 工程师与设计师联合撰写专业细节
  • 统一编辑规范,保证文风与排版一致
  • 每月复盘转化数据,迭代选题方向

延伸阅读

相关资讯与近期热门内容

深度阅读推荐

建站决策前值得细读的三篇

网站改版的5个关键决策
2024-08-12

网站改版的5个关键决策

什么时候该改版、改到什么程度、如何避免流量掉光,京华建材集团改版复盘给出答案。

获取专属建站方案

看完文章,把您的行业与预算告诉我们,免费获取一份量身定制的官网建设方案与报价。

立即免费咨询