Apache Spark Tungsten 与 AQE 深度实战:从 UnsafeRow 内存布局、Whole-Stage Code Generation 到自适应查询执行的工程全解
关键词:Spark / Tungsten / AQE / Whole-Stage CodeGen / UnsafeRow / 倾斜 Join / ColumnarBatch
很多人对 Spark 调优的认知停留在"调并行度、加 executor 内存、开 AQE"这三板斧上。但真正决定一个 Job 是跑 8 分钟还是 40 分钟的,往往是两件被埋藏很深的东西:数据在内存里到底长什么样(Tungsten),以及执行计划能不能在运行时变(AQE)。这篇文章就沿着 Spark SQL 的物理执行栈往下挖,把这两层的工程实现和落地陷阱讲透。
一、为什么要 Tungsten:JVM 对象模型的原罪
Spark 早期版本中,一行 Product、Order、User 构成的 JOIN 结果会被表达成一个个 Java 对象树。以 (Int, String, Long) 三列的行举例,JVM 里的真实开销是:
- 每个对象 16 字节对象头(Object Header,含 Mark Word 与 Klass Pointer);
String内部还有一个char[](Java 8)或byte[](Java 9+ 紧凑字符串)引用,再加一层 object header、length 字段,实际 4 字节字符可能吃掉 40+ 字节;- GC 需要遍历整条对象引用链,百万行 JOIN 意味着百万棵对象树要标记、复制、晋升。
结果是:有效数据占比可能不到 10%,剩下的全是元数据和 GC 压力。Tungsten 项目的目标非常直接——把内存布局从「JVM 对象图」拉回到「C 结构体」。
1.1 UnsafeRow:一行数据的扁平二进制形式
UnsafeRow 是 Tungsten 的核心。它的内存布局严格分成三段:
┌────────────────┬────────────────────┬─────────────────────────────┐
│ null bit set │ fixed-length │ variable-length data │
│ 8 bytes/64bit │ values area │ (strings / arrays / maps) │
└────────────────┴────────────────────┴─────────────────────────────┘
设计要点非常值得学习:
numFields不存。行里只放 bitset 和数据,Schema由算子持有并传递。这意味着你无法单独解释一个UnsafeRow,必须配 Schema——这是用「schema 外置」换来的每行 4 字节节省。- 8 字节对齐。每个字段从 offset = bitsetSize + 8 * i 处读取,允许 JVM 以对齐的长字访问。
- 变长字段两阶段写——变长字段(String/Array/Map)先写入尾部 variable-length 区,头部区固定留 8 字节存放 (offset, size) 对。这样定长区永远可以直接按
ordinal * 8随机寻址,无需扫描。
读一个字段的代码路径大致是这样(逻辑简化版):
// UnsafeRow.java 的核心读取逻辑(示意)
long offsetInBits = getLong(baseObject, baseOffset); // null bitmap
long relativeOffset = baseOffset + bitSetWidthInBytes; // 跳过 bitset
if (isNullAt(ordinal)) return null; // 查 bitmap
long valueOffset = relativeOffset + ordinal * 8L;
return Platform.getLong(baseObject, valueOffset); // sun.misc.Unsafe 直接读取
注意这里用的是 sun.misc.Unsafe(Spark 3 之后通过 Platform 抽象层包装,并处理 JDK17+ 的 --add-opens 强封装问题),做的是裸地址偏移读取,没有虚方法分派、没有边界检查、没有对象解引用。这就是数据中心里常说的「SIMD 的另一个世界:这里靠的不是向量化指令,而是把 random pointer chase 变成 linear scan」。
1.2 Off-Heap 还是 On-Heap?
Spark 提供两种内存模式:
spark.memory.offHeap.enabled = true
spark.memory.offHeap.size = 8g
生产建议:默认不要开 off-heap。 理由往往被忽略:
UnsafeRow在 heap 内模式一样是紧凑二进制布局,已经解决了 90% 的痛点;- off-heap 模式下
Platform.getLong读取 byte[] 之外的内容,JIT 无法把它优化成单条指令,反而可能变慢; - off-heap 内存不再被 GC 追踪,泄露不会自动回收,一旦出现
TaskMemoryManager泄漏,只能靠重启 executor。
真正适合 off-heap 的场景只有两个:堆大 (>32G) 想躲 ZGC/Region 开销,或者和堆外 IO 库(如某些 native sort / arrow 池)共存。
二、Whole-Stage Code Generation:让 Java 跑得像手写 C
2.1 火山模型的代价
经典 Volcano / Iterator 模型里,SELECT ... FROM a JOIN b WHERE ... GROUP BY 需要挨个调用 hasNext() / next()。每个算子接口调用都是虚方法分派,CPU 需要:查 vtable、做分支预测、把 Row 对象从栈上移下去。更糟的是中间结果被反复装箱/拆箱,CPU 流水线几乎全在 stall。
2.2 思路:把算子「熔」,生成一段 Java 源码
Spark 的做法是——既然我们本来就在 JVM 上,那就动态生成一段 Java 代码,把 Scan → Filter → Project → Aggregate 全部融进一个 while 循环,然后用 Janino 编译器现场编译成一个类。
生成出来的代码结构大致是下面这样(为可读性做了简化):
// 生成的 WholeStageCodegenExec 的内部循环模板
while (scan_input.hasNext()) {
InternalRow row = (InternalRow) scan_input.next();
// 1) Filter: a.value > 100 AND b.type = 'PAID'
boolean filterPassed = false;
long v = row.getLong(2);
UTF8String t = row.getUTF8String(5);
if (!(false || row.isNullAt(2))) {
filterPassed = v > 100L && t.equals(PAID_STR);
}
if (!filterPassed) continue;
// 2) Project: key为 concat(a.id, '_', b.channel)
UTF8String key = UTF8String.concat(row.getUTF8String(0), SEP, row.getUTF8String(7));
// 3) Aggregate: 写入 UnsafeFixedWidthAggregationMap
long idx = hashMap.findOrInsert(key);
agg_buf[idx] += v; // 内联的聚合累加,无对象分配
}
关键收益:
- 中间对象完全不再产生 ——
InternalRow自始至终只有一份,零 GC 压力; - 分支预测友好 —— 循环体里所有字段访问都是确定偏移量;
- JIT 能进一步把
Platform.getLong优化成极少数指令;
但也触发两个隐藏成本:
- Janino 编译延迟。复杂 SQL 生成的方法体可能超过 JVM 64KB bytecode 上限,触发 Split。
spark.sql.codegen.wholeStage=true(默认开)时,一个 40 层的 case-when / 超宽聚合会退化为多个 stage。 - 方法体膨胀导致 JIT 放弃。单次 size 极大的循环体可能使 C2 编译器拒绝内联,反而比火山模型更慢。这时应手动:
spark.sql.codegen.maxFields // 默认 200,超过就拆分对象为多个方法
spark.sql.codegen.hugeMethodLimit // 默认 65536 bytes
诊断方式很直接:看 EXPLAIN CODEGEN。它会逐个算子标注是否支持代码生成;没有 WholeStageCodegen 前缀包裹的算子,说明它触发了 fallback(不支持 codegen 或超过了上述门槛)。
2.3 ShuffleExchangeExec:为什么 Exchange 是 Codegen 的天然边界
ShuffleExchangeExec 必须被 tense 到 mark presence 缺口:Shuffle 是全局操作,不可融合,所以它天然划分 codegen stage。这就是 Spark UI 里「Stage ≠ codegen stage」的由来:一个 SQL Stage 内部可能由多个 codegen pipeline 组成,中间被 sort / aggregate(final) / exchange 打断。
三、AQE:在运行时重写你自己
编译期优化器(Catalyst)依赖统计信息做代价估算,但真实世界里:
- 中间结果与行数无从预估(多杈 JOIN、UDF、Python UDF);
- 数据倾斜使
avg row size完全没有意义; - Table stats 是 stale 的(分区堆半小时前的)。
AQE(Adaptive Query Execution)的思路是:先跑起来,跑完一小段就重新规划。ShuffleExchangeExec 是天然的执行边界:它必须先把自己 doExecute() 完并上报 MapOutputStatistics,下游才能开始 shuffle read。这个「先写后读」的窗口,正好给 AQE 留下了重新规划的时机。
AQE 的三大核心能力:
3.1 动态合并 Shuffle 分区(CoalesceShufflePartitions)
静态 spark.sql.shuffle.partitions=2000 是万恶之源:数据量小的时候 2000 个分区产生大量小文件和调度开销,数据量大时又不够。AQE 会在每个 shuffle read 完成后,按实际大小把相邻 partition 当 Small Parts 粘并:
spark.sql.adaptive.enabled = true
spark.sql.adaptive.coalescePartitions.enabled = true
spark.sql.adaptive.advisoryPartitionSizeInBytes = 256MB // 目标尺寸,别太小
spark.sql.adaptive.coalescePartitions.minPartitionNum = 1
实战观点:advisoryPartitionSizeInBytes 默认 64MB 偏小。容器内单个 task 最佳处理量常在 128–256MB 之间(要看你的列宽和 spill 阈值)。不要盲目往下调,分区过碎会让调度开销和 task 序列化成本吃掉全部收益。
3.2 动态切换 Join 策略(DemoteBroadcastHashJoin / OptimizeLocalShuffleReader)
Catalyst 可能因为 stats 不准,把一个实际只有 20MB 的表判成了 SortMergeJoin。AQE 在 shuffle write 完成后,可以从 shufflewrite_statuses 中算出真实的 map side sizes,一旦发现它其实小于:
spark.sql.autoBroadcastJoinThreshold = 64MB // 可配到 200MB 以内
就把执行计划里的 SortMergeJoin 动态换为 BroadcastHashJoin。用户侧的 EXPLAIN 看不出来(plan 是 rewritten),要到 Spark UI 的 SQL tab 里看底部的 alternative plan。
3.3 倾斜 Join 自动拆分(OptimizeSkewedJoin)
这是 AQE 实用性最强的一条。当一个 reducer 的输入远超中位数时:
spark.sql.adaptive.skewJoin.enabled = true
spark.sql.adaptive.skewJoin.skewedPartitionFactor = 5 // 中位数倍数
spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes = 256MB
Spark 会把超大 reduce partition 拆成 N 个小任务,在 map 侧也把对应的 shuffle 数据复制 N 份,然后再 UNION 结果。这就是用空间换时间:拖尾的那一个 task 从读 10GB 变成 N 个 task 各读 10GB/N。
生产坑:这套机制只对 SMJ(SortMergeJoin) 生效。如果你的 JOIN 是 BroadcastHashJoin,倾斜处理完全不触发。另一个坑是 skewedPartitionFactor=5 太粗,对本身分布就很宽的 stage(如 group by 后的长尾)会产生大量误判,白白引入 map 侧数据复制的开销——建议结合历史 metrics 调整到 10 以上。
四、'hello world' 到生产:一个完整的调优闭环
import org.apache.spark.sql.functions._
val df = spark.read.parquet("s3://dw/fact_orders/dt=2026-*")
.filter(col("dt") >= "2026-01-01")
.join(broadcast(dimUser), Seq("user_id"), "left")
.join(dimMerchant.hint("MERGE"), Seq("mch_id"), "left")
.groupBy("mch_id", "channel")
.agg(sum("amount").as("gmv"))
df.write.mode("overwrite").parquet("s3://dw/agg_daily/")
诊断三步法:
- 先看对位置:不要看预估 plan,看 Spark UI 的
SQL / DataFrametab,滚动到最下面查看 AQE 实际执行后的分层计划。 - 再验证 codegen:用
EXPLAIN formatted观察每个算子是否被WholeStageCodegen包裹。 - 最后才调内存:
Execution Memory与Storage Memory在 UnifiedMemoryManager 下是软边界,可以互相借用。一旦看到 spill,往往暗示 partition 太小、或 agg key 基数超预期。
五、结论
Tungsten 和 AQE 分别解决了 Spark 性能问题的两个半边:
- Tungsten 解决的是「每行数据搬多少字节」——用紧凑二进制布局 + fused codegen 把JVM的抽象成本压到近似 C 的水平;
- AQE 解决的是「计划能不能变」——把优化决策从「编译期猜测」迁移到「运行时实测」。
这也是现代计算引擎的共同趋势:数据布局决定下限,自适应执行决定上限。Velox、Flink Adaptive Batch Scheduler、Trino FTE 都在做同一件事。理解这两层的实现细节,比背十个参数有用得多。
真正专业的做法不是记住 spark.sql.xxx=yyy,而是在 EXPLAIN 出执行计划之后,能用 Platform.getLong 的思维去理解每一行代码里 CPU 到底在做什么。

发表评论 取消回复