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)  │
└────────────────┴────────────────────┴─────────────────────────────┘

设计要点非常值得学习:

  1. numFields 不存。行里只放 bitset 和数据,Schema 由算子持有并传递。这意味着你无法单独解释一个 UnsafeRow,必须配 Schema——这是用「schema 外置」换来的每行 4 字节节省。
  2. 8 字节对齐。每个字段从 offset = bitsetSize + 8 * i 处读取,允许 JVM 以对齐的长字访问。
  3. 变长字段两阶段写——变长字段(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 优化成极少数指令;

但也触发两个隐藏成本:

  1. Janino 编译延迟。复杂 SQL 生成的方法体可能超过 JVM 64KB bytecode 上限,触发 Split。spark.sql.codegen.wholeStage=true(默认开)时,一个 40 层的 case-when / 超宽聚合会退化为多个 stage。
  2. 方法体膨胀导致 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/")

诊断三步法:

  1. 先看对位置:不要看预估 plan,看 Spark UI 的 SQL / DataFrame tab,滚动到最下面查看 AQE 实际执行后的分层计划。
  2. 再验证 codegen:用 EXPLAIN formatted 观察每个算子是否被 WholeStageCodegen 包裹。
  3. 最后才调内存: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 到底在做什么。

点赞(0) 打赏

评论列表 共有 0 条评论

暂无评论
立即
投稿

微信公众账号

微信扫一扫加关注

发表
评论
返回
顶部