Spark 硬核指南(中篇):引擎盖下的 Spark,与它撑起的生态
承接上篇。如果说上篇是"这辆车能带你去哪、和其他车比起来怎么样",中篇就是打开引擎盖——看看 Catalyst 优化器怎么把你的代码变成执行计划,Tungsten 怎么榨干 CPU 和内存,Shuffle 到底在传什么,以及 Spark 之上撑起的整个 Lakehouse 生态。这一篇会比上篇更"硬核",建议配合 Spark UI 实际操作一遍,理解会深得多。
1. Catalyst 优化器:从 DataFrame 代码到执行计划的四道关卡
上篇留了一个悬念:为什么 DataFrame/SQL 写法天生比手写 RDD 快?答案就是 Catalyst——Spark SQL 引擎的查询优化器,本质上是一个基于 Scala 函数式编程特性构建的规则驱动 + 代价驱动的查询优化框架,思路上和传统数据库(Oracle、PostgreSQL)的优化器一脉相承,但用现代函数式语言重新实现,扩展性更强。
1.1 第一关:解析(Parsing)
无论你写的是 SQL 字符串还是链式调用的 DataFrame API,最终都会被解析成一棵未绑定(Unresolved)的逻辑计划树——这一步只管语法对不对,完全不检查表名、列名是否真实存在。
1.2 第二关:分析(Analysis)
结合 Catalog(元数据目录,记录了表结构、列类型、UDF 注册信息)把"未绑定"的树,逐个节点绑定成真实的表、列、函数引用。这一步如果你写错了列名,就是在这里报 AnalysisException。
1.3 第三关:逻辑优化(Logical Optimization)——规则驱动
这是最容易在面试中被问到细节的一步。Catalyst 会应用一整套启发式规则,与具体的物理执行方式无关,纯粹在逻辑层面做等价变换:
| 优化规则 | 作用 | 直观例子 |
|---|---|---|
| 谓词下推(Predicate Pushdown) | 把 WHERE 条件尽量往数据源方向推 | filter 从 join 之后挪到 join 之前,甚至挪进 Parquet 文件的 Reader 里,跳过不满足条件的行组 |
| 列裁剪(Column Pruning) | 只保留后续真正用到的列 | SELECT a FROM t 只读 Parquet 文件里 a 这一列,其余列的字节完全不读 |
| 常量折叠(Constant Folding) | 编译期算出常量表达式的值 | WHERE 1 + 1 = 2 直接被替换成 WHERE true |
| 布尔表达式简化 | 化简冗余的逻辑表达式 | a OR true 简化为 true |
| Join 重排序 | 调整多表 Join 的顺序 | 结合表的统计信息,让小表先参与 Join,减少中间结果 |
这些规则是**穷举式地反复应用,直到树不再变化(定点迭代)**为止。
1.4 第四关:物理计划生成——代价驱动
逻辑优化只回答"做什么",物理计划要回答"怎么做"——尤其是 Join 策略的选择,这是性能差异最大的一环:
| Join 策略 | 触发条件 | 特点 |
|---|---|---|
| Broadcast Hash Join | 一侧表小于广播阈值(默认 10MB,可调) | 把小表完整广播到每个 Executor,避免 Shuffle,最快 |
| Shuffle Hash Join | 两侧都较大,但一侧明显小于另一侧 | 双侧按 Join Key Shuffle 后在内存建哈希表 |
| Sort Merge Join | 两侧都很大,是默认兜底策略 | 双侧按 Key 排序后归并,稳定但涉及 Shuffle+排序,开销最大 |
这里先埋一个伏笔:Spark 2.x 时代,Join 策略是在物理计划生成阶段"一次性拍板"的,如果表的实际大小和统计信息(Statistics)不准(比如经过复杂过滤后行数远小于预估),就可能选错策略,跑出很慢的 Sort Merge Join。这个问题在 Spark 3.0 引入的 AQE 里被彻底解决——本篇第 3 节详细讲。
1.5 用 explain() 亲眼看穿一切
不用猜,直接问引擎:
df4.explain(mode="formatted")
# 或者看更详细的四个阶段
df4.explain(True) # 会依次打印 Parsed/Analyzed/Optimized/Physical 四棵树给你的建议:任何一次"这段代码为什么这么慢"的排查,第一步永远是 explain(),不是瞎猜。看物理计划里有没有出现意料之外的 SortMergeJoin、有没有 Exchange(Shuffle 的物理算子名)、Scan 节点有没有把你以为已经下推的 filter 体现出来。这是比看 Spark UI 更早、更便宜的排查手段。
2. Tungsten 执行引擎:从"JVM 拖后腿"到"接近手写代码"
Catalyst 把逻辑规划做到了极致,但如果物理执行层还是传统 JVM 对象套对象的方式,性能天花板依然很低——JVM 对象有巨大的内存开销(对象头、指针、装箱开销),垃圾回收(GC)会在数据量大的时候造成剧烈的停顿,Java 序列化又慢又占空间。Tungsten 项目就是奔着这三座大山去的。
2.1 堆外内存 + 二进制行格式
Tungsten 不再用 Java 对象表示一行数据,而是用自定义的紧凑二进制格式(UnsafeRow),直接操作堆外内存(Off-Heap Memory)里的字节数组。好处是三重的:
- 省内存:没有对象头、没有指针追踪开销,同样的数据占用空间显著减少;
- 躲开 GC:堆外内存不受 JVM 垃圾回收管理,数据量再大也不会造成 GC 停顿飙升;
- CPU Cache 友好:紧凑的二进制布局对 CPU Cache Line 更友好,减少 Cache Miss。
2.2 全阶段代码生成(Whole-Stage CodeGen)
这是 Tungsten 里最"黑科技"的部分。传统的火山模型(Volcano Model)执行引擎,每个算子(filter/project/aggregate)都是一个独立的迭代器对象,数据一行一行地在算子之间通过虚函数调用传递——每传一行数据,就是一次虚函数调用,CPU 分支预测频繁失败,性能损耗巨大。
Whole-Stage CodeGen 的做法是:把一个 Stage 内一连串的算子(比如 filter → project → hash aggregate 前半段),在运行时直接生成成一段专用的 Java 字节码,编译执行,效果上接近你自己手写一个 for 循环处理所有逻辑——彻底消灭了算子间的虚函数调用开销。
# 用这个方法能看到 Spark 实际生成的 Java 代码(了解即可,不需要真的读懂每一行)
df4.explain("codegen")一句话记住 Tungsten 的核心思路:与其让通用的解释执行引擎慢慢跑,不如在运行时"现写现编译"一段专用代码——这正是"用代价换灵活性"这条老路的反向操作:"用规划开销换执行速度"。
3. AQE:自适应查询执行,运行时的第二次机会
前面埋的伏笔在这里兑现。AQE(Adaptive Query Execution,自适应查询执行)是 Spark 3.0 引入、目前默认开启的重大特性,核心思想是:静态优化基于估算的统计信息,估算会出错;不如让引擎在运行过程中,用真实产生的中间结果统计信息,重新规划后续步骤。
AQE 主要解决三类问题:
3.1 动态合并 Shuffle 分区(Coalescing Post-Shuffle Partitions)
静态阶段规划的 Shuffle 分区数(默认 200)往往是"一刀切"的经验值,实际运行后某些分区可能只有几 KB 数据。AQE 会在 Shuffle 完成后,看实际每个分区的数据量,自动合并过小的相邻分区,减少后续 Task 数量和调度开销。
3.2 动态切换 Join 策略
前面说过,物理计划阶段选错 Join 策略是老大难问题。AQE 允许在 Shuffle 完成、拿到真实数据量后,如果发现某一侧数据其实很小,动态把原本规划好的 Sort Merge Join 改写成 Broadcast Hash Join——这在多层过滤之后的复杂查询里,效果非常显著,往往是几倍到十几倍的性能提升。
3.3 动态处理数据倾斜的 Join(Skew Join Optimization)
这是生产环境里最救命的一个特性。数据倾斜(某个 Join Key 对应的数据量远超其他 Key,比如电商场景里"未知用户"这个 Key 占了 40% 的订单)会导致处理这个 Key 的 Task 严重拖后腿,整个 Stage 卡在最慢的那个 Task 上。AQE 能检测到某个分区明显大于其他分区的中位数,自动把这个过大的分区拆成多个子分区分别参与 Join,再合并结果,从根源上缓解倾斜(下篇会讲手动处理倾斜的加盐等经典技巧,AQE 让很多场景不再需要手动介入)。
# AQE 相关的核心配置(4.x 默认已开启,了解参数含义即可)
spark.conf.set("spark.sql.adaptive.enabled", "true") # 总开关,默认 true
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true") # 动态合并分区
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true") # 倾斜 Join 优化
spark.conf.set("spark.sql.adaptive.advisoryPartitionSizeInBytes", "128m") # 建议分区大小给你的建议:如果你还在维护 Spark 2.x 的老代码库,升级到 3.x/4.x 拿到 AQE 免费的性能提升,往往比手动调优半天参数更划算——这是过去几年 Spark 版本升级里投入产出比最高的一次架构升级,没有之一。
4. Shuffle 内幕:分布式计算里最贵的一步
如果说前三节讲的是"引擎怎么变聪明",这一节讲的是"分布式计算里那件最贵、最容易出问题的事到底是什么"。
Shuffle 的本质:把数据按某个 Key 重新分布到不同分区(比如 groupBy、join、repartition),这意味着数据要跨网络、跨节点重新洗牌——涉及磁盘写、网络传输、磁盘读三重开销,是 Spark 作业里最昂贵的操作,没有之一。
4.1 Shuffle 的物理过程
Spark 目前使用的是 Sort-based Shuffle(早期还有 Hash-based Shuffle,已被淘汰):Map 端每个 Task 把输出数据按 Key 排序后写成一个文件(配合一个索引文件标记每个分区的偏移量),Reduce 端的 Task 按需去对应位置**拉取(fetch)**自己该处理的数据段。
4.2 Push-based Shuffle:新一代 Shuffle 优化
传统 Pull 模式下,Reduce 端要主动去成百上千个 Map Task 的输出位置逐个拉取小数据块,在大规模集群上会产生大量随机小 I/O、连接数暴增的问题。Push-based Shuffle(Spark 3.2 引入,需要额外部署 Shuffle Service,社区里 Apache Celeborn 等项目是这个方向更进一步的独立中间件方案)改变思路:Map 端主动把输出推送、合并到少量的 Shuffle 服务节点上,Reduce 端只需要从少数几个大文件里顺序读取,把大量随机小 I/O 变成少量顺序大 I/O,在超大规模集群上收益明显。
4.3 一句话总结 Shuffle 的调优哲学
能不 Shuffle 就不 Shuffle(用 Broadcast Join 代替 Shuffle Join);必须 Shuffle 时尽量减少数据量(先 filter/select 再 join,而不是反过来);Shuffle 分区数要匹配数据量和集群规模(这是下篇性能调优的核心内容之一)。
5. 内存管理:统一内存模型
Spark 的 Executor 内存不是"一大坨随便用",而是被划分成几个区域,理解这张图,是排查 OOM(下篇会有真实案例)的地基:
统一内存管理(Unified Memory Management,Spark 1.6 之后的默认模型)的核心创新是:Execution 内存和 Storage 内存不再是死板的固定比例,而是可以互相"借用"——如果当前没有缓存任务在跑,Execution 可以用满整个统一内存区做 Shuffle/排序;如果 Storage 需要缓存空间而 Execution 正在使用,会按照一套优先级规则(Execution 优先于 Storage,可以强制驱逐部分缓存)动态调整。这比 1.6 之前"两块内存死磕不能互相借用"的旧模型灵活得多。
6. 调度系统:DAGScheduler 与 TaskScheduler
回到上篇的 Application → Job → Stage → Task 层级,这一节讲清楚"谁负责把这个层级关系真正构建出来"。
- DAGScheduler:面向 Stage 层面的调度器。它拿到一个 Job 的完整 RDD 依赖链后,按 Shuffle(宽依赖)边界切分出 Stage(这正是配套笔记里"宽依赖划分新 Stage"这条规则的具体执行者),并维护 Stage 之间的依赖关系(DAG);
- TaskScheduler:面向 Task 层面的调度器。它接过 DAGScheduler 给的一批 Task,**结合数据本地性(Locality)**把 Task 分发到合适的 Executor 上执行——优先把 Task 调度到数据所在节点,尽量避免跨节点读取数据;
- 推测执行(Speculative Execution):如果一个 Task 的执行时间明显长于同 Stage 其他 Task(可能是硬件问题、数据倾斜等原因),调度器可以在其他节点推测性地启动一个相同的 Task 副本,谁先跑完就用谁的结果,另一个被杀掉——这是应对"长尾任务"的经典手段,代价是会消耗额外的集群资源,需要权衡开启。
# 推测执行相关配置
spark.conf.set("spark.speculation", "true")
spark.conf.set("spark.speculation.multiplier", "1.5") # 比中位数慢 1.5 倍就触发推测7. Structured Streaming 深入:微批模型的工程实现
上篇对比过 Spark 和 Flink 的流处理哲学差异,这里深入 Structured Streaming 具体怎么实现"看起来像流"的批处理。
7.1 核心心智模型:无界表(Unbounded Table)
Structured Streaming 的编程模型极其优雅:把持续到达的流数据,当成一张不断有新行追加进来的"无界表",你写的查询逻辑和批处理 DataFrame 代码几乎一模一样,引擎在背后每隔一个触发间隔(Trigger Interval)就对新增的数据跑一次增量计算。
# 从 Kafka 读取流数据,做窗口聚合,几乎和批处理代码一样
stream_df = (spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "kafka1:9092")
.option("subscribe", "orders")
.load())
result = (stream_df
.selectExpr("CAST(value AS STRING) as json")
.select(from_json(col("json"), schema).alias("data"))
.select("data.*")
.withWatermark("event_time", "10 minutes")
.groupBy(window(col("event_time"), "5 minutes"), col("product_id"))
.agg(sum("amount").alias("total_amount")))
query = (result.writeStream
.format("console")
.outputMode("update")
.trigger(processingTime="10 seconds") # 微批触发间隔
.start())7.2 Watermark:处理"迟到数据"的核心机制
流数据天然是乱序的——用户手机断网 5 分钟后事件才补发上来,这类"迟到数据"该怎么处理?Watermark(水位线)是引擎对"我们能容忍多迟的数据"划的一条线:withWatermark("event_time", "10 minutes") 表示"最多接受比当前观测到的最大事件时间晚 10 分钟到达的数据,之后这个时间窗口就彻底关闭、清理状态,更晚到的数据会被丢弃"。这是在结果准确性和状态无限增长之间做的权衡——不设 Watermark,状态会无限累积,最终撑爆内存。
7.3 Exactly-once 语义是怎么做到的
Structured Streaming 声称端到端 Exactly-once(精确一次),靠的是三重保障:数据源可重放(Kafka 的 offset 机制)+ Checkpoint 记录精确的处理进度(WAL 预写日志)+ Sink 端幂等写入或事务写入。三者缺一不可——如果你的 Sink 是一个不支持幂等写入的自定义系统,即便 Spark 端做得再好,端到端也达不到 Exactly-once,这是实践中最容易被忽视的一个坑。
7.4 Trigger 模式一览
| Trigger 模式 | 行为 | 延迟特性 |
|---|---|---|
| 默认(未指定) | 上一批处理完立刻开始下一批 | 尽可能快,但仍是微批 |
processingTime | 固定间隔触发一次微批 | 可预测的延迟,最常用 |
once | 只处理一次当前所有可用数据后停止 | 适合"半小时跑一次"这种伪流场景 |
availableNow(3.3+) | 处理完当前所有可用数据(可能分多个微批)后停止 | 比 once 更高效地处理积压数据,替代旧的 once |
continuous(实验性) | 尝试逐行处理,追求毫秒级延迟 | 支持的算子有限,生产成熟度不及微批模式,目前仍非主流选择 |
给你的建议:continuous 模式听起来很诱人(毫秒级延迟),但支持的算子集合有限、生产环境案例稀少。如果业务真的需要毫秒级延迟,更现实的选择是直接换 Flink,而不是赌 Spark 的实验性特性——这也是上篇"选型建议"背后的深层原因。
8. MLlib 与 GraphX:不是本系列重点,但要知道它们存在
- MLlib:Spark 的机器学习库,提供了特征工程(
VectorAssembler、StringIndexer等)、经典算法(逻辑回归、随机森林、梯度提升树、K-Means 等)、以及贯穿全流程的PipelineAPI(把特征处理和模型训练串成一条可复用、可持久化的流水线)。它的强项是"大规模结构化数据上的经典机器学习",不擅长深度学习——深度学习训练目前的主流方案是 Spark 做数据预处理、Ray/Horovod/原生分布式 PyTorch 做训练,中间用 Parquet/Delta 或 Petastorm 这类桥接层衔接。 - GraphX / GraphFrames:图计算库,处理社交网络分析、推荐系统的图算法(PageRank、连通分量等)。GraphX 是基于 RDD 的老 API,GraphFrames 是社区维护的基于 DataFrame 的替代方案,性能和易用性都更好,新项目应优先考虑 GraphFrames。
9. Lakehouse 生态全景:Spark 之上的存储层战争
如果说前面几节讲的是 Spark 引擎本身,这一节要讲清楚为什么 2020 年代后,几乎所有新建的数据平台都在说"我们是 Lakehouse 架构",以及 Spark 在其中扮演什么角色。
9.1 问题的起点:数据湖不是数据库
传统数据湖(一堆 Parquet/CSV 文件躺在 HDFS/S3 上)没有事务保证、没有 Schema 演进管理、不支持高效的更新删除、也没有"时间旅行"(回溯历史版本)能力——这些恰恰是数据库最基本的能力。**开放表格式(Open Table Format)**要解决的就是这个问题:在文件之上加一层元数据管理,让数据湖具备接近数据库的 ACID 事务、Schema 演进、Time Travel 能力,同时保留数据湖"存储计算分离、廉价对象存储、多引擎可读"的优势。
目前主流的三大开放表格式:
| 格式 | 诞生背景 | 核心设计哲学 |
|---|---|---|
| Delta Lake | Databricks(2019) | Spark 原生优先,与 Databricks Runtime 深度绑定 |
| Apache Iceberg | Netflix(2018) | 规范先行(Spec-first),任意引擎都能实现读写,多引擎中立性是核心卖点 |
| Apache Hudi | Uber(2016) | 为高频流式 Upsert/CDC 场景而生,内置增量处理和表管理服务 |
三种格式诞生的业务背景至今仍深深烙印在各自的架构里:Hudi 生于 Uber 对海量实时车程事件做 Upsert 的需求,Iceberg 生于 Netflix 在 S3 上做 PB 级分析、同时想摆脱 Hive 表的种种限制,Delta Lake 生于 Databricks 想在 Spark 工作负载上获得 ACID 保证。
9.2 2026 年的格局:Iceberg 领跑,但不是"赢家通吃"
2026 年,Iceberg 凭借厂商中立的治理模式、分区演进能力,以及横跨 Spark、Flink、Trino、Snowflake、BigQuery、DuckDB 的最广泛多引擎支持,已经成为事实上的行业标准;Hudi 在纯流式/CDC 摄入场景仍保持领先;Delta Lake 在以 Databricks 为中心的环境里依然强势。Apache Iceberg v3 规范已于 2026 年年中达到正式发布状态,补齐了删除向量(Deletion Vector)、行血缘(Row Lineage)、VARIANT 类型等此前相对其他格式的功能差距;与此同时,各大云厂商和数据平台都已支持读写 Iceberg——包括 Delta Lake 的创造者 Databricks 自己。
Apache Hudi 在 2026 年迎来了 1.0 这一里程碑式的重构:采用 LSM 结构化时间线、非阻塞并发控制(让写入、Compaction、Clustering 互不阻塞地并行推进)、以及作为一等公民的二级索引子系统。这让 Hudi 在处理高频更新/删除的场景(比如订单状态频繁变更、用户资料修改)上,依然保有独特优势——"合并读取(Merge-on-Read)"表类型允许更新只追加小的日志文件,而不必重写整个 Parquet 文件,后台异步 Compaction 再慢慢合并,这个设计对写多读少的 CDC 场景是降维打击式的优化。
9.3 该怎么选?一张决策流程图
一个值得记住的行业真相是:查询引擎本身的性能往往比表格式的选择更影响最终体验——不管是 Databricks 上的 Photon、Trino 配合 Iceberg,还是一个调优良好的 Spark 集群,三种主流格式在大多数场景下都能跑出不错的结果,实际项目里更值得投入精力做针对性基准测试,而不是纠结"哪个格式绝对最快"。
9.4 用 Spark 读写 Iceberg 表:一个实际的例子
spark = (SparkSession.builder
.appName("iceberg-demo")
.config("spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions")
.config("spark.sql.catalog.local", "org.apache.iceberg.spark.SparkCatalog")
.config("spark.sql.catalog.local.type", "hadoop")
.config("spark.sql.catalog.local.warehouse", "s3://my-bucket/warehouse")
.getOrCreate())
# 建表、写入,语法和普通 Spark SQL 表几乎一致
spark.sql("""
CREATE TABLE local.db.orders (
order_id BIGINT, user_id BIGINT, amount DOUBLE, event_time TIMESTAMP
) USING iceberg
PARTITIONED BY (days(event_time))
""")
df.writeTo("local.db.orders").append()
# 时间旅行:查询某个历史快照
spark.sql("SELECT * FROM local.db.orders VERSION AS OF 123456789").show()
# Schema 演进:加列不需要重写全表数据
spark.sql("ALTER TABLE local.db.orders ADD COLUMN discount DOUBLE")10. 部署模式再深入:Standalone / YARN / Kubernetes 该怎么选
上篇只是速览,这里给出更贴近生产决策的对比:
| 维度 | Standalone | YARN | Kubernetes |
|---|---|---|---|
| 部署复杂度 | 最简单,Spark 自带 | 需要完整 Hadoop 生态 | 需要 K8s 集群,云托管服务可省去大部分运维 |
| 资源隔离 | 较弱 | 成熟的队列/资源池机制 | 容器级隔离,配合 Namespace/ResourceQuota 精细控制 |
| 弹性伸缩 | 较弱 | 中等 | 原生支持,与云厂商 Auto Scaling 无缝集成 |
| 多租户支持 | 弱 | 强(企业级资源调度成熟) | 强,且和其他云原生工作负载统一调度 |
| 与 Lakehouse/云存储集成 | 需自行配置 | 需自行配置 | 云原生工具链(如 Spark Operator)集成度高 |
| 典型场景 | 自建小规模集群、快速验证 | 已有 Hadoop 存量集群的企业 | 新建集群、云原生优先的团队 |
给你的建议:如果你是从零开始搭建一个新的数据平台(没有历史 Hadoop 集群包袱),2026 年的默认答案基本是 Kubernetes——原因不只是"技术更先进",更现实的是运维团队的技能栈正在整体向云原生迁移,用同一套 K8s 技能同时管理 Spark、其他微服务、AI 训练任务,比维护一套独立的 YARN 技能栈成本更低。如果你已经有成熟的 YARN 集群和团队经验,没有必要为了"追新"而强行迁移。
11. Spark Connect:客户端与服务端的解耦
这是近几个版本里一个容易被低估、但工程价值很大的架构变化。传统 Spark 应用里,Driver 和你的客户端代码是同一个进程——这意味着你的 Notebook/IDE 崩了,Driver 也跟着没了;你想用一个轻量的 Python 环境远程操作一个庞大的生产集群,也很别扭。
Spark Connect(3.4 引入,4.0 起大幅增强) 把架构改成了标准的客户端-服务端模式:客户端只需要一个极轻量的库(不需要完整的 Spark 发行版),把 DataFrame 操作序列化成协议缓冲区(基于 gRPC)发给远程的 Spark Connect Server,服务端才是真正跑着完整 Spark 引擎的地方——这和数据库客户端连接远程数据库服务器的模式几乎一模一样,是 Spark 架构上一次"补课式"的现代化。
这个改动带来的实际好处:客户端崩溃不影响集群作业、可以用一个极轻量的 Python 客户端库连接远程集群(不需要在本地装完整 JVM 环境)、多语言客户端更容易实现(Python、Java 之外,Go、Rust、Swift 等语言客户端的实现难度大幅降低,因为只需要实现协议层,不需要移植整个执行引擎)。
# Spark Connect 客户端写法:只是把 builder 换成 remote 连接
from pyspark.sql import SparkSession
spark = SparkSession.builder.remote("sc://my-spark-cluster:15002").getOrCreate()
df = spark.read.parquet("s3://bucket/data.parquet")
df.groupBy("category").count().show()
# 你的 Python 进程里没有真正跑 Spark 引擎,所有计算都发生在远程集群上本篇小结
- Catalyst 用四道关卡(解析→分析→逻辑优化→物理计划)把你的代码变成执行计划,规则驱动 + 代价驱动,
explain()是排查一切"为什么这么慢"的第一步; - Tungsten 用堆外内存二进制格式和全阶段代码生成,把 JVM 的性能短板补了回来,本质是"用规划开销换执行速度";
- AQE 是过去几年投入产出比最高的架构升级,让引擎能根据运行时真实统计信息重新决策 Join 策略、合并小分区、处理数据倾斜;
- Shuffle 是分布式计算里最贵的操作,理解它的物理过程,是后续所有性能调优的地基;
- 统一内存管理让 Execution 和 Storage 内存能动态互相借用,理解这张图能帮你更快定位 OOM 根因;
- Structured Streaming 用"无界表"心智模型统一了批流编程体验,但本质仍是微批,Watermark 和 Exactly-once 语义的实现细节决定了它在什么场景够用、什么场景不够用;
- Lakehouse 生态(Iceberg/Delta/Hudi)是 Spark 之上正在发生的最重要的架构演进,2026 年 Iceberg 是多引擎场景的默认选择,Delta 在 Databricks 生态内依然强势,Hudi 在高频 CDC 场景保有独特优势;
- Spark Connect 是架构现代化的重要一步,把 Driver 和客户端解耦,为多语言客户端和更轻量的开发体验铺路。
下篇预告:从原理回到实战——性能调优的具体手法(分区策略、Join 优化、数据倾斜的手动处理技巧、序列化选择)、生产环境真实踩过的坑(OOM、小文件、GC 调优)、Spark UI 怎么读、成本优化的思路、Spark 4.x 的最新特性、GPU 加速与 Photon 引擎的对比,以及面试和职业发展的实用建议。