
前阵子帮朋友调一个跑批任务凌晨两点被叫起来说是生产上一个大SQL跑了四个小时还没完集群明明有几百个核磁盘I/O却一直在报警。我上去一看Spark UI第一个Stage的Shuffle Write把数据全落到了磁盘之后每个Stage都在等I/OCPU根本没吃满。这种场景在大数据领域太典型了——数据仓库跑批变慢、交互查询卡顿十有八九不是算力不够而是内存没用对。这篇就聊聊我在数仓场景里反复打磨过的几类内存优化策略从执行引擎参数、数据组织、缓存复用到倾斜治理完整讲一遍思路希望能给同样在扛数仓的同学一些能直接上手的参考。1. 瓶颈不在算力在搬运重新理解数仓查询慢的根源先说一个反直觉的结论大多数数仓任务跑得慢不是CPU不够而是数据的搬运成本太高。1.1 磁盘与内存的速度差距到底有多大我做优化之前习惯先算一笔带宽账。机械硬盘顺序读写的速度大概在150MB/s到200MB/sSATA固态在500MB/s左右NVMe固态能到2GB/s以上听着已经很快了。但内存呢一台普通的双通道服务器理论带宽轻松到几十GB/s高配多通道机器上百GB/s也很正常。也就是说内存和磁盘之间隔着至少一到两个数量级的差距。如果你的任务有一半时间花在等磁盘I/O上那一味加CPU核数、加并发度基本是在拿金子补墙。数据仓库恰恰是重I/O场景。数仓跑批链路里从HDFS扫描原始数据开始到Shuffle洗牌、Join、Aggregate中间结果一旦超过内存容量就必须落盘。落盘一次读盘一次本来几秒能算完的东西被拖到几分钟甚至几十分钟。所以我认为内存优化最底层的一句话是减少数据在内存和磁盘之间的搬运次数提高单位内存的产出效率。1.2 数仓链路里的每一次落盘都在烧时间传统Hive数仓用MapReduce跑SQL时这个特征最明显。一个简单的JoinMap端把数据写本地磁盘Reduce端再从磁盘拉数据每个Stage之间都要经历一次写盘-读盘的循环。Spark虽然有内存计算的优势但一旦spill内存溢出到磁盘发生照样会陷入同样的困境。所以我在做优化时心里始终有一个清单按开销从大到小排列任务级落盘Shuffle Write、Shuffle Read中间结果全量读写磁盘应用级落盘不合理的persist缓存级别、RDD重复计算后写中间表数据级低效全表扫描、读入大量根本用不到的列、高压缩比的存储格式没选对资源级浪费Executor堆内存配得很大但Execution和Storage区域互相抢占、GC频繁内存成了摆设后面的所有策略其实都在围绕这张清单展开。2. 内存优化从执行引擎开始Spark和Hive的参数怎么调才有效参数调优是整个内存优化里最容易被误解的部分。不少人一上来就把executor内存加到8G、16G结果任务该慢还是慢甚至更慢。原因是他们没搞懂Spark的内存管理模式。2.1 Spark统一内存管理模型拆解Spark从2.x开始采用统一内存管理Unified Memory Management。一个Executor的JVM堆内内存大致被划分成保留区域Reserved默认300MB、可用执行内存Execution、可用存储内存Storage和其他用户内存User Memory。其中Execution和Storage是重点它们共享同一个内存池比例由两个参数控制spark.memory.fractionExecutor堆内可用于Execution和Storage管理的比例默认0.6剩下0.4留给用户代码、内部元数据和保护性空间。spark.memory.storageFractionStorage在共享内存池中优先占用的比例默认0.5。这段设计有个很聪明的点存储内存缓存RDD/DataFrame和Execution内存Shuffle、Join、Aggregate临时数据可以互相借用。缓存的数据在需要时可以被Execution强制驱逐驱逐不出去就写磁盘。也就是说内存不是物理隔离的而是按压力动态挤占。也正是因为这种挤占机制如果你的任务里缓存了大量数据同时又在跑大Shuffle两边就会互相抢内存最后缓存被写盘、Shuffle也在spill整体性能直接雪崩。这个问题很多人遇到过但不知道该怪谁其实就是storageFraction没调好。2.2 生产环境可落地的参数组合我在不同场景下做过几组组合目前比较稳的配置逻辑是这样的场景关键参数推荐值说明OLAP交互查询spark.executor.memory4G - 8G查询并发多单任务内存不宜过大OLAP交互查询spark.memory.fraction0.6 - 0.7预留空间给用户代码和任务调度OLAP交互查询spark.memory.storageFraction0.3 - 0.4交互查询临时计算多Storage让位ETL批量跑批spark.executor.memory8G - 16G单任务吞吐优先ETL批量跑批spark.memory.fraction0.7 - 0.75计算密集尽量让内存服务计算ETL批量跑批spark.memory.storageFraction0.3跑批很少需要长时间缓存多留给Execution混部集群spark.executor.memoryOverheadmax(executorMem * 0.1, 384m)堆外内存不足会报native memory不足所有场景spark.sql.adaptive.enabledtrueSpark 3.x AQE强烈建议打开重点说下ETL场景。我的习惯是适当调大spark.memory.fraction但不要硬调spark.memory.storageFraction超过0.5。因为跑批任务里最常见的操作是GroupBy、Join这些吃的是Execution内存。如果把Storage占得太多等于把最紧缺的资源优先给了不常用的缓存Shuffle一膨胀就只能spill。另一个容易被忽视的是spark.executor.memoryOverhead。它负责JVM之外的堆外内存包括线程栈、网络缓冲、Native库。数据量大时如果这个值不够会出现Container killed on request. Exit code is 143或java.lang.OutOfMemoryError: Direct buffer memory这类报错。很多人以为堆内存不够拼命加executor内存其实问题出在堆外。2.3 Join策略也得跟着内存走Spark在做Join时有个很关键的参数spark.sql.autoBroadcastJoinThreshold默认值是10MB。意思是当某张小表小于这个阈值时Spark会把它广播到每个Executor内存里避免Shuffle。这个参数能在很大程度上改变查询对内存的压力。比如一张3GB的维度表在跑大量事实表Join维度表时如果不广播每一次Join都要做全量Shuffle几亿条数据的洗牌代价非常大。广播之后每个Executor只需要在内存里放一份3GB的维度数据Shuffle直接消失。但广播不是免费的。你要清楚广播总量是表大小×Executor数量叠加的。如果Executor有200个、小表10GB理论上要占2TB内存这谁也扛不住。所以我的经验是广播策略适合几百MB到一两GB的维表超过5GB或Executor特别多的集群要谨慎。Hive侧对应的是hive.auto.convert.join.noconditionaltask.size默认值一般10MB左右作用类似。在Hive数仓里做MapJoin时调大这个阈值能显著减少Reduce阶段的Shuffle数据量。要注意Hive的MapJoin是小表加载到Map端内存阈值太大同样有内存风险我一般控制在100MB-200MB。3. 数据进门之前先瘦身存储格式、压缩与分区裁剪执行引擎参数调完下一步不是继续加内存而是回到源头让更少的数据被读进内存。3.1 列式存储为什么天然适合内存优化数仓里最常见的低效操作是对着一个文本文件或行式存储的表跑全表扫描把所有列全部读进内存哪怕SQL里只用了3个字段。列式存储ORC、Parquet解决的就是这个问题。它的核心思路是按列存放数据查询时只读取涉及的列再配合谓词下推和列剪枝能大幅压缩扫描量。真实验证中同一份TPC-DS数据集从TextFile换成ORC扫描数据量经常能少70%-80%。这意味着内存里驻留的数据少了后续Join、聚合碰到的数据量也就小了整个链路的压力都在下降。在Hive/Spark数仓里我强烈建议把ODS之上所有层级的表统一成ORC格式建表类似这样CREATE TABLE dws_order_daily ( order_date STRING, province_id INT, order_cnt BIGINT, amount DOUBLE ) PARTITIONED BY (dt STRING) STORED AS ORC TBLPROPERTIES (orc.compress ZSTD);当初我们组把核心宽表从Parquet切到ORCZSTD时有几个任务直接快了一倍多就是因为Parquet在应对细粒度行级更新和Metadata读取时开销偏大ORC在Hive生态下配合度更好。3.2 压缩算法选型不能只看压缩比压缩算法是内存优化的第二层筛子。压缩比高意味着同样大小的表读取时产生的临时数据更少内存占用更小但压缩和解压要消耗CPU选型要平衡。算法压缩比压缩/解压速度适合场景Gzip高慢归档数据极少查询的冷表Snappy中快日常查询和ETL兼容性最好的默认选择ZSTD高较快追求压缩比和速度均衡的数仓热表LZ4低极快对延迟极敏感的实时链路我现在的默认组合是ORC ZSTD。ZSTD在压缩比上接近Gzip速度却快得多尤其在CPU核数充足的集群上非常划算。不过要注意如果你的集群CPU数很紧张、I/O带宽又富余那Snappy可能更合适压缩比低一点但省CPU。3.3 分区、分桶与文件治理都是隐形内存优化分区裁剪的作用是跳过无关数据。数仓里必须按常用过滤维度做好分区设计比如日期、地区、业务线。我见过一个表没有分区每次跑批全量扫描3TB数据加了日期分区后单次任务只扫当天80GB内存压力瞬间消失。分桶则能优化Join的Shuffle策略。当两张表在Join字段上分桶数量一致时Spark/Hive可以做Bucket Join同一个桶的数据在本地关联不需要全量Shuffle。分桶数一般按数据量和Executor并行度来定我常用的是32、64、128这类与Executor数量成比例的数值。再就是小文件治理。一个分区下塞了几万个小文件光打开文件、拿元数据就要占大量Driver和Executor内存扫数据时I/O也碎得没法看。我的习惯是跑批结束后通过spark.sql.adaptive.coalescePartitions.enabled自动合并小分区或者定期用INSERT OVERWRITE ... SELECT重刷一遍表把碎片文件压缩到合理数量。4. 缓存与复用让同一份数据只进一次内存参数调好了数据也减重了下一个目的是同样一份数据能不能不要反复从磁盘读4.1 哪些数据值得缓存一个简单的判断标准我在指导组里新人时经常说一句话缓存是给高频、中等大小、低变更的数据准备的。高频同一个DataFrame或Hive表在一个任务里或连续多个任务里被反复用到。中等大小总体积在集群内存可承受范围内一般几十MB到几GB。低变更维表、配置表、日活UV这种相对稳定的数据。反过来几TB的事实表、每天都在全量变化的结果表都不适合长缓存。事实表塞进缓存大概率把整个Executor内存打爆还会跟Shuffle抢空间得不偿失。4.2 persist级别选择与缓存淘汰Spark里cache()其实等价于persist(StorageLevel.MEMORY_ONLY)。更精细的控制要看persistStorageLevel描述适用场景MEMORY_ONLY只放内存放不下则重新计算RDD派生成本低数据量小MEMORY_AND_DISK内存放不下就落盘不重新计算最常用安全与性能的折中MEMORY_AND_DISK_SER序列化后放内存内存占用更小数据较大且GC压力高DISK_ONLY全落盘几乎不推荐缓存意义降低我最常用的组合是MEMORY_AND_DISK_SER结合Kryo序列化。Kryo比Java默认序列化节省将近一半的内存占用特别适合在缓存大量K-V结构数据时用。配置方式spark.conf.set(spark.serializer, org.apache.spark.serializer.KryoSerializer)缓存也不是存完就完了。任务结束后不主动释放缓存缓存对象会一直占着Executor内存拖垮后续任务。我现在习惯在DataFrame使用结束后显式调用unpersist()或者在开发规范里要求每个Spark任务结束时清理自己的缓存。4.3 从缓存到结果复用物化视图和结果表单任务内的缓存只能解决一个应用生命周期内的重复计算。跨任务级的复用靠的是把中间结果落地成表。Hive的典型做法是分层建模ODS/DWD/DWS每层算完写表下游直接读表而不是一路嵌套子查询重算。这本质上是把计算结果而不是原始数据放进存储减少下游扫描量。Hive 3还支持物化视图Materialized View查询引擎可以自动改写SQL命中物化视图对重复的聚合口径省掉重新计算。在Spark侧连接池或结果表也能实现类似效果——用一个写好的结果表替代每次从源头开始跑。这个方向经常被忽视但我认为它比任何参数调优都省内存。因为无论Execution内存怎么调永远比不上根本不读那些数据来得彻底。5. 容易被忽略的隐性内存杀手数据倾斜、Shuffle与GC真实生产环境里参数和数据格式都合理任务还是慢往往是因为隐性杀手在作怪。5.1 数据倾斜一个热点Key压垮整个Executor数据倾斜是数仓优化里最经典的内存杀手。表现也很典型某个Stage里绝大多数Task几十秒跑完但有一两个Task要跑半小时甚至直接OOM或者某个Executor内存爆掉整个Application被反复重试。原因通常是group by的key分布不均一个热点值比如某个城市、某个用户的占比特别大把几亿条数据塞进了同一个分区。整个Executor的Execution内存被这个Task吃光剩下的数据全部spill到磁盘于是越跑越慢最后GC和磁盘I/O双重压力下崩溃。处理思路有两类加盐两阶段聚合先给key加随机前缀将热点值打散到多个分区做局部聚合再去掉前缀做全局聚合。适合Aggregate场景。交给AQE自动优化Spark 3.x开启spark.sql.adaptive.skewJoin.enabled后动态检测倾斜的Shuffle分区并按比例拆分能自动缓解部分倾斜问题。AQE不是银弹复杂业务还是得自己在SQL层面加盐。但把开关打开至少能兜住一部分偶发倾斜。5.2 Shuffle溢出与Tungsten堆外内存Shuffle的spill机制号称不会OOM因为放不下就写磁盘。但这意味着任务的实际耗时从秒级变成分钟级。内存优化里的一个重要工作就是看Spark UI里的Spill指标一旦发现Spill (Memory)和Spill (Disk)有数值就要警惕了。Tungsten是Spark的底层优化模块它用Unsafe内存管理直接操作堆外内存避免JVM对象开销。合理设置spark.memory.offHeap.enabled和spark.memory.offHeap.size把部分Shuffle缓冲放到堆外能够减少GC压力。不过堆外内存的上限受系统物理内存约束设置过高反而会挤占页面缓存让操作系统频繁换页。我的经验是shuffle量特别大比如几十GB以上的任务可以试着让5%-10%的执行内存走堆外同时把spark.shuffle.file.buffer从默认32KB适当调大减少Shuffle写的小I/O次数。5.3 GC长暂停加内存反而变慢的元凶很多同学说内存越大越好我不这么认为。JVM堆超过32GBGC停顿会明显拉长尤其在对象创建极快的Shuffle场景。Spark UI里每个Executor的GC Time如果超过总运行时间的10%就要小心了。我的做法是分两步确认序列化方式有没有换Kryo没换先换这是性价比最高的降GC方式。如果GC仍然严重检查是否用了G1GC并在spark-submit里加上-XX:UseG1GC -XX:MaxGCHPauseMillis200。对于大堆场景G1比Parallel GC的停顿控制更好。但也要明白GC调优是锦上添花真正的药方还是减少数据进出内存的量而不是换一个GC器。6. 一套可以抄作业的优化流程从定位到验证前面都是单一策略最后我把整套优化流程串起来。这也是我每次处理数仓慢任务的固定动作。6.1 先看Spark UI三个指标快速定位内存瓶颈拿到一个慢任务我不会先改参数而是先打开Spark UI做三轮检查Stage列表页找出执行时间最长、Shuffle Read数据量最大的Stage。这个大Stage才是优化的主战场。进到Stage详情看Shuffle Read Size / Records如果单Task读取量有几十GB说明Shuffle量过大优先考虑Join策略和过滤下推。看Spill指标一旦看到Memory/Disk Spill有值任务已经在向磁盘求救这个Stage的内存压力是重点。还有一个容易忽略的看Executor页的Memory和GC Time。如果GC时间占比高说明JVM堆内对象的存活量太大需要从序列化、缓存上控内存。6.2 小步验证改一项、测一项、对比基线我的习惯是建立一套基线数据用固定的几条代表SQL一个复杂Join、一个GroupBy聚合、一个维表关联做基准测试。每次只改一个变量跑3次取中位数记录执行时间和资源占用。比如优化项优化前优化后提升表存储TextFile → ORCZSTD18分20秒9分05秒约50%开启BroadcastJoin小表广播9分05秒4分32秒约50%加盐处理热点Key4分32秒2分48秒约38%调整memoryFraction和storageFraction2分48秒2分22秒约15%能看到越往后的优化收益越小。这也是我坚持先减数据量、再调执行策略、最后调内存参数的原因——数据量减下来后面参数的收益才会放大。6.3 把经验沉淀成开发规范优化从来不是一个人的事。我调完一次任务后会把结论固化到数仓开发规范里。比如建表统一ORCZSTD所有ETL任务尽量命中分区裁剪Join时小表优先广播禁止对超大宽表执行笛卡尔积式关联对group by热点字段做预聚合或加盐处理每段Spark代码结束时主动unpersist禁止缓存泄漏跑批任务的核心SQL统一纳入基线性能测试这套规范执行半年后我们组里数仓核心跑批的整体耗时下降了一大截而且新任务的返工率明显降低。很多问题在开发阶段就被拦住了。最后再说一点个人体会。内存优化这项工作表面看是调参本质上是让数据在正确的位置上流动。我踩过的坑多了之后最大的心得是别一上来就堆内存先看数据长什么样再决定内存怎么用。跑批前花两分钟写条探针SQLselect count(*) from table where dt ...顺带分字段看一眼唯一值分布很多倾斜、膨胀问题在动手前就能暴露出来比事后盯Spark UI要省太多事。