Spark 完全入门指南(中篇):技术栈与生态对比
上篇讲了 Spark 的核心原理,这篇讲 Spark 的五大组件,以及和其他大数据技术的深度对比。
面向有经验的工程师,不只是"是什么",更讲"为什么这么设计"、"和其他技术比怎么样"、"什么场景该用什么"。
目录
- 一、Spark SQL:不只是"用 SQL 查大数据"
- 二、Structured Streaming:新一代流计算引擎
- 三、MLlib:分布式机器学习的定位与边界
- 四、GraphX:图计算的现状与替代方案
- 五、Spark vs Flink:大数据领域的"世纪之争"
- 六、Spark vs MapReduce / Storm / Tez:代际对比
- 七、Spark SQL vs Hive vs Presto/Trino vs ClickHouse:SQL 引擎选型
- 八、生态集成:Spark 与 Hive/Kafka/HBase/Delta Lake 的协作
一、Spark SQL:不只是"用 SQL 查大数据"
1.1 Spark SQL 的定位:从"SQL 接口"到"统一 API"
很多人以为 Spark SQL 就是"让你能用 SQL 查 Spark 数据",这只是表面。
实际上,Spark SQL 是 Spark 2.0 之后的核心 API 层:
- DataFrame / Dataset API 构建在 Spark SQL 之上
- Structured Streaming 构建在 Spark SQL 之上
- MLlib 的新 API(spark.ml)构建在 DataFrame 之上
- 甚至 GraphX 的新替代品 GraphFrames 也构建在 DataFrame 之上
💡 关键洞察: Spark 2.0 之后,RDD 虽然还在,但已经变成了"底层实现"。 上层应用开发,官方推荐用 DataFrame/Dataset/SQL,因为有 Catalyst 优化器和 Tungsten 执行引擎的加持,性能比手写 RDD 更好。 这是一个重要的架构转变:从"程序员手写优化"到"优化器自动优化"。
1.2 DataFrame vs Dataset vs RDD:到底用哪个?
| 特性 | RDD | DataFrame | Dataset |
|---|---|---|---|
| 类型安全 | ✅ 编译时检查 | ❌ 运行时才发现 | ✅ 编译时检查 |
| 优化器 | ❌ 无 | ✅ Catalyst | ✅ Catalyst |
| 执行引擎 | 通用 | ✅ Tungsten | ✅ Tungsten |
| 语言支持 | Scala/Java/Python/R | 全支持 | 仅 Scala/Java |
| 性能 | 基准 | 比 RDD 快 2~10 倍 | 和 DataFrame 相当 |
| 灵活性 | 最高 | 中等 | 中等 |
| 适用场景 | 底层操作、自定义计算 | 大多数场景、SQL 风格 | 强类型需求、复杂业务逻辑 |
选型建议
| 场景 | 推荐 |
|---|---|
| 数据清洗、ETL、SQL 查询 | DataFrame / SQL |
| 复杂业务逻辑、需要类型安全 | Dataset(Scala/Java) |
| 底层操作、自定义分区、特殊计算 | RDD |
| Python 用户 | DataFrame(Python 没有 Dataset) |
🎯 一句话总结: 90% 的场景用 DataFrame/SQL 就够了,性能好、写起来简单。 RDD 只在需要底层控制的时候用。 Dataset 是 Scala/Java 用户的福利,Python 用户用不了。
1.3 Catalyst 优化器:Spark SQL 快的秘密
Catalyst 是 Spark SQL 的查询优化器,基于 Scala 的函数式编程特性构建,支持规则优化和代价优化。
查询执行的四个阶段
SQL / DataFrame 代码
│
▼
┌─────────────────┐
│ 1. 解析(Analysis) │ → 把 SQL 解析成未解析的逻辑计划,解析表名、列名、类型
└────────┬────────┘
▼
┌─────────────────┐
│ 2. 逻辑优化(Logical Optimization)│ → 规则优化:谓词下推、列裁剪、常量折叠...
└────────┬────────┘
▼
┌─────────────────┐
│ 3. 物理计划(Physical Planning)│ → 生成多个物理计划,用代价模型选最优
└────────┬────────┘
▼
┌─────────────────┐
│ 4. 代码生成(Code Generation)│ → 生成 Java 字节码, Whole-Stage CodeGen
└────────┬────────┘
▼
执行查询常见的优化规则
| 优化规则 | 作用 | 例子 |
|---|---|---|
| 谓词下推(Predicate Pushdown) | 把过滤条件尽量往下推,减少数据量 | 先 filter 再 join,而不是先 join 再 filter |
| 列裁剪(Column Pruning) | 只读取需要的列,不读不需要的 | select 只查 3 列,就只读 3 列 |
| 常量折叠(Constant Folding) | 把常量表达式提前算出来 | where 1+1=2 直接变成 where true |
| 投影合并(Projection Merge) | 合并相邻的投影操作 | 多个 map 合并成一个 |
| Join 策略选择 | 根据数据大小选择合适的 Join 算法 | 小表用 Broadcast Join,大表用 Sort Merge Join |
| 空值传播(Null Propagation) | 利用空值语义简化表达式 | null + x 直接是 null |
💡 为什么 Catalyst 很牛? 它意味着你写的 SQL 哪怕不怎么优化,Catalyst 也会帮你优化成比较高效的执行计划。 这和传统数据库的查询优化器是一个思路,但 Spark 的 Catalyst 是开源的、可扩展的,你可以自己加优化规则。
怎么看执行计划?
// 看逻辑计划和物理计划
df.explain(true)
// 只看物理计划
df.explain()看执行计划是调优 Spark SQL 的第一步,一定要学会看。
1.4 Tungsten 执行引擎:榨干 CPU 性能
Tungsten(钨丝计划)是 Spark 1.5 引入的执行引擎优化,目标是榨干 CPU 和内存的性能。
三大核心优化
1. 内存管理和二进制处理
- 直接操作二进制数据,不创建 JVM 对象
- 用
sun.misc.Unsafe直接操作 off-heap 内存 - 减少 JVM 对象开销(对象头、对齐填充)和 GC 压力
传统方式:数据 → JVM 对象(有对象头、指针、对齐)→ 操作 → GC 回收 Tungsten:数据 → 二进制(紧凑存储)→ 直接操作 → 几乎无 GC
2. 缓存感知计算(Cache-aware computation)
- 算法设计考虑 CPU 缓存命中率
- 数据按缓存行大小组织,减少缓存失效
- 类似数据库的火山模型(Volcano Model)优化
3. 全阶段代码生成(Whole-Stage Code Generation)
- 把整个查询阶段的算子融合成一个 Java 方法
- 消除虚函数调用、中间数据物化
- 生成的字节码接近手写优化代码
传统的火山模型(算子之间每条数据都要调用 next()):
Scan → Filter → Project → HashAggregate
↑ ↑ ↑ ↑
每条数据在算子之间传递,有大量函数调用开销
Whole-Stage CodeGen(把整个阶段编译成一个方法):
for (每条数据) {
if (过滤条件) {
计算投影
聚合
}
}
没有函数调用,数据在寄存器/缓存里流转🎯 性能提升有多大? 官方数据,Tungsten 让 Spark SQL 的性能提升了 2~10 倍。 特别是 Whole-Stage CodeGen,是 Spark SQL 比 RDD 快的主要原因之一。
1.5 Adaptive Query Execution(AQE):运行时动态优化
Spark 3.0 引入的 AQE 是一个大杀器,解决了静态优化器的痛点。
静态优化器的问题
Catalyst 是静态优化,在执行前就定好执行计划。但问题是:
- 统计信息可能不准(比如表的行数、列的分布)
- 数据倾斜在编译时发现不了
- Join 的数据量实际运行时才知道
结果就是:优化器选的计划可能不是最优的。
AQE 怎么解决?
AQE 在运行时根据实际数据统计信息,动态调整执行计划。
三大核心功能
| 功能 | 作用 | 效果 |
|---|---|---|
| 动态合并 Shuffle 分区 | 运行时把小的分区合并,减少 Task 数 | 避免小文件、小 Task 过多 |
| 动态切换 Join 策略 | 运行时发现小表,自动从 SortMergeJoin 换成 BroadcastJoin | Join 性能大幅提升 |
| 动态优化数据倾斜 | 运行时发现倾斜的分区,自动拆分 | 解决数据倾斜导致的长尾 Task |
怎么开启?
# Spark 3.0+ 默认是关闭的,3.2+ 默认开启
spark.sql.adaptive.enabled=true
spark.sql.adaptive.coalescePartitions.enabled=true
spark.sql.adaptive.skewJoin.enabled=true🎯 AQE 的意义: 以前调优要手动算分区数、手动处理倾斜、手动判断要不要 broadcast。 开了 AQE 之后,很多调优工作 Spark 自动帮你做了。 这是 Spark 从"程序员手动调优"向"自调优"演进的重要一步。
二、Structured Streaming:新一代流计算引擎
2.1 从 Spark Streaming 到 Structured Streaming
Spark Streaming 的问题
Spark Streaming(DStream)是 Spark 1.x 的流计算方案,基于微批处理。它的问题:
- API 不统一:DStream API 和 DataFrame API 是两套,批处理和流处理代码不能复用
- 延迟较高:最小批处理间隔 0.5 秒,实际一般 1~5 秒
- 事件时间支持弱:只能用处理时间,不支持基于事件时间的窗口
- 状态管理难:有状态计算要自己实现,容错复杂
- 不支持 exactly-once:只能 at-least-once
Structured Streaming 的改进
Spark 2.0 推出 Structured Streaming,基于 Spark SQL,核心思想:
把流数据当成一张"无限增长的表",你对这张表写查询,就像批处理一样。
批处理:有限的表 → 查询 → 有限的结果
流处理:无限增长的表 → 持续查询 → 持续更新的结果这就是"批流统一"的思路:同一套 API、同一个优化器、同一种执行引擎。
2.2 核心概念
输入表(Input Table)
流数据源产生的数据,不断追加到这张无限表里。
结果表(Result Table)
查询的结果,每次触发计算都会更新。
输出(Output)
结果表写到外部存储的方式,有三种输出模式:
| 输出模式 | 说明 | 适用场景 |
|---|---|---|
| Append(默认) | 只输出新增的行 | 简单的 ETL、没有聚合 |
| Complete | 每次输出全部结果 | 有聚合、需要全量结果 |
| Update | 只输出有变化的行 | 有聚合、只关心变化 |
触发器(Trigger)
多久计算一次:
| 触发方式 | 说明 |
|---|---|
| 固定间隔(ProcessingTime) | 每隔 N 秒触发一次 |
| 一次性(OneTime) | 只跑一次,适合定时任务 |
| 连续处理(Continuous) | 真正的流处理,毫秒级延迟(实验性) |
| 可用时立即(AvailableNow) | 有数据就处理,处理完等下一批(Spark 3.3+) |
2.3 事件时间与 Watermark
这是流计算的核心难点之一。
什么是事件时间?
- 处理时间(Processing Time):数据到达 Spark 的时间
- 事件时间(Event Time):数据本身携带的时间(比如日志产生的时间)
为什么需要事件时间?因为数据可能乱序到达、延迟到达。 比如用户 10:00 产生的日志,可能因为网络延迟 10:05 才到 Spark。 如果按处理时间统计,就会统计错误。
Watermark(水位线)
Watermark 定义了"允许数据延迟多久"。
import org.apache.spark.sql.functions._
val windowedCounts = df
.withWatermark("eventTime", "10 minutes") // 允许延迟 10 分钟
.groupBy(
window($"eventTime", "5 minutes") // 5 分钟窗口
)
.count()工作机制:
- Watermark = 当前看到的最大事件时间 - 延迟阈值
- 窗口结束时间 < Watermark 的窗口,就关闭,不再接受延迟数据
- 关闭的窗口结果可以输出(Append 模式)
💡 为什么需要 Watermark? 如果没有 Watermark,Spark 要永远记住所有窗口的状态(因为不知道数据会不会还来),内存会爆。 Watermark 告诉 Spark:超过这个时间的延迟数据我不要了,你可以清理状态了。
2.4 有状态计算与状态管理
Structured Streaming 内置支持有状态计算:
- 窗口聚合
- 去重(dropDuplicates)
- 流流 Join(Stream-Stream Join)
- 自定义有状态处理(mapGroupsWithState / flatMapGroupsWithState)
状态存储后端:
- HDFS 状态存储(默认,适合小状态)
- RocksDB 状态存储(适合大状态,Spark 2.3+)
⚠️ 状态管理是流计算的核心难点: 状态太大 → OOM 状态恢复慢 → 重启时间长 状态过期策略不合理 → 内存泄漏或数据丢失 生产环境一定要监控状态大小,合理设置 Watermark。
2.5 Structured Streaming vs Spark Streaming
| 对比项 | Spark Streaming (DStream) | Structured Streaming |
|---|---|---|
| 基于 | RDD | DataFrame/Dataset |
| API | DStream 专用 API | 和批处理统一 |
| 延迟 | 秒级(最小 0.5s) | 秒级(可到 100ms) |
| 事件时间 | 不支持 | 原生支持 + Watermark |
| 状态管理 | 需自己实现 | 内置支持 |
| 语义保证 | at-least-once | exactly-once(部分 sink) |
| 推荐度 | 维护模式 | 官方推荐 |
| 学习曲线 | 较低 | 中等(事件时间、状态等概念) |
🎯 结论: 新项目直接用 Structured Streaming,不要再用 DStream 了。 DStream 已经处于维护模式,不会有大的功能更新。 老项目如果还在用 DStream,可以考虑迁移。
三、MLlib:分布式机器学习的定位与边界
3.1 MLlib 的定位:不是"分布式 TensorFlow"
很多人听到 Spark MLlib 就以为是"分布式深度学习框架",这是误解。
MLlib 的定位是:分布式传统机器学习库。
- ✅ 擅长:传统机器学习(逻辑回归、随机森林、K-Means、协同过滤等)
- ✅ 擅长:数据已经在 Spark 生态里,不想搬来搬去
- ✅ 擅长:特征工程(TF-IDF、标准化、特征选择等)
- ❌ 不擅长:深度学习(神经网络、CNN、RNN、Transformer)
- ❌ 不擅长:超大规模模型训练(那是参数服务器/分布式训练框架的事)
💡 关键认知: MLlib 和 TensorFlow/PyTorch 不是竞争关系,而是互补关系。
- 数据预处理、特征工程 → 用 Spark(数据量大)
- 模型训练(深度学习)→ 用 TensorFlow/PyTorch(GPU 加速)
- 模型推理/批量预测 → 可以用 Spark 批量跑(如果模型能导出)
3.2 两套 API:spark.mllib vs spark.ml
| API | 基于 | 状态 | 推荐 |
|---|---|---|---|
| spark.mllib | RDD | 维护模式 | 不推荐 |
| spark.ml | DataFrame | 活跃开发 | 官方推荐 |
spark.ml 的核心概念是 Pipeline(工作流):
DataFrame → Transformer → Transformer → Estimator → Model → 预测- Transformer:把一个 DataFrame 变成另一个 DataFrame(比如特征转换)
- Estimator:训练数据,产生一个 Transformer(比如算法训练出模型)
- Pipeline:把多个 Transformer 和 Estimator 串起来
import org.apache.spark.ml.Pipeline
import org.apache.spark.ml.feature.{Tokenizer, HashingTF}
import org.apache.spark.ml.classification.LogisticRegression
// 定义流水线
val tokenizer = new Tokenizer().setInputCol("text").setOutputCol("words")
val hashingTF = new HashingTF().setInputCol("words").setOutputCol("features")
val lr = new LogisticRegression()
val pipeline = new Pipeline().setStages(Array(tokenizer, hashingTF, lr))
// 训练(自动执行所有阶段)
val model = pipeline.fit(trainingData)
// 预测
val predictions = model.transform(testData)💡 Pipeline 的好处:
- 特征工程和模型训练串在一起,不会漏步骤
- 可以一起保存/加载,部署方便
- 可以做交叉验证和参数搜索(CrossValidator + ParamGridBuilder)
3.3 MLlib 包含什么?
特征工程(最常用)
| 类别 | 算法 |
|---|---|
| 文本特征 | TF-IDF、Word2Vec、CountVectorizer |
| 数值转换 | StandardScaler、MinMaxScaler、MaxAbsScaler、Normalizer |
| 特征选择 | ChiSqSelector、VectorSlicer |
| 特征构造 | PolynomialExpansion、Interaction、Bucketizer |
| 编码 | StringIndexer、OneHotEncoder、IndexToString |
| 缺失值 | Imputer |
分类
| 算法 | 特点 |
|---|---|
| 逻辑回归 | 二分类/多分类,线性模型,可解释性强 |
| 决策树 | 可解释性强,容易过拟合 |
| 随机森林 | 集成学习,泛化能力强,常用 |
| 梯度提升树(GBT) | 集成学习,效果好,训练慢 |
| 朴素贝叶斯 | 文本分类常用,简单快速 |
| 多层感知机(MLP) | 简单的神经网络,仅支持浅层 |
| 线性 SVM | 二分类,线性模型 |
回归
| 算法 | 特点 |
|---|---|
| 线性回归 | 基础回归模型 |
| 决策树回归 | 非线性回归 |
| 随机森林回归 | 集成回归,常用 |
| 梯度提升树回归 | 效果好,训练慢 |
| 生存回归 | 生存分析用 |
| Isotonic Regression | 保序回归 |
聚类
| 算法 | 特点 |
|---|---|
| K-Means | 最常用,需要指定 K |
| 高斯混合(GMM) | 软聚类,每个点属于多个簇的概率 |
| LDA | 主题模型,文本分析 |
| Bisecting K-Means | 二分 K-Means,更快 |
推荐
| 算法 | 特点 |
|---|---|
| ALS(交替最小二乘) | 协同过滤,显式/隐式评分都支持 |
频繁模式挖掘
| 算法 | 特点 |
|---|---|
| FP-Growth | 频繁项集挖掘 |
| PrefixSpan | 序列模式挖掘 |
3.4 MLlib 的局限与替代方案
局限
- 算法更新慢:很多新算法(XGBoost、LightGBM、CatBoost)没有原生支持
- 不支持深度学习:MLP 只是简单的前馈网络,没有 CNN/RNN/Transformer
- 超参搜索慢:网格搜索是串行的,没有高级的超参优化
- 模型可解释性工具少:没有 SHAP、LIME 等工具
- GPU 支持弱:主要靠 CPU,GPU 加速还不成熟
替代/补充方案
| 需求 | 推荐方案 |
|---|---|
| 梯度提升树 | XGBoost4J / LightGBM(有 Spark 集成) |
| 深度学习 | TensorFlow / PyTorch(数据用 Spark 预处理) |
| 超参优化 | Hyperopt / Optuna(可以和 Spark 集成) |
| 模型可解释性 | SHAP(可以用 Spark 批量计算) |
| AutoML | 各家 AutoML 平台 |
🎯 生产环境的常见架构: 数据湖 → Spark 做特征工程 → 特征存到特征存储 → TensorFlow/PyTorch 训练模型 → 模型部署到推理服务
Spark 在这个流程里负责"数据准备"和"批量预测",模型训练交给专业的深度学习框架。
四、GraphX:图计算的现状与替代方案
4.1 GraphX 是什么?
GraphX 是 Spark 的图计算组件,处理图数据(顶点 + 边)。
核心抽象:
- VertexRDD:顶点 RDD,(VertexId, VD)
- EdgeRDD:边 RDD,Edge(srcId, dstId, ED)
- Graph:图,由 VertexRDD + EdgeRDD 组成
内置算法:
- PageRank
- 连通分量
- 三角形计数
- 最短路径
- 社区发现(Louvain 没有内置,需要自己实现或用 GraphFrames)
4.2 GraphX 的现状
说实话,GraphX 现在比较尴尬:
- 维护不活跃:GraphX 很久没有大的功能更新了
- API 基于 RDD:没有迁移到 DataFrame,享受不到 Catalyst/Tungsten 的优化
- 算法少:只有基础算法,很多高级算法没有
- 性能一般:和专门的图计算引擎比有差距
4.3 替代方案
| 方案 | 类型 | 特点 | 适用场景 |
|---|---|---|---|
| GraphFrames | Spark 上的图库 | 基于 DataFrame,API 更现代,算法更多 | 已经在用 Spark,需要图计算 |
| Neo4j | 图数据库 | 最流行的图数据库,查询语言 Cypher 强大 | 图数据为主,需要频繁查询 |
| JanusGraph | 分布式图数据库 | 基于 HBase/Cassandra,可扩展 | 超大规模图数据 |
| DGL / PyG | 图神经网络框架 | 深度学习 + 图,做 GNN | 图神经网络、节点预测 |
| Flink Gelly | Flink 的图组件 | 和 GraphX 类似,也不太活跃 | 已经在用 Flink |
🎯 建议:
- 如果只是偶尔做简单的图计算(比如 PageRank),用 GraphX 或 GraphFrames 就行
- 如果图是核心业务,用专门的图数据库(Neo4j/JanusGraph)
- 如果要做图神经网络,用 DGL/PyG
- GraphX 了解基本概念就行,不用深入研究
五、Spark vs Flink:大数据领域的"世纪之争"
这是大数据领域最常被问到的对比,也是最容易产生"信仰之争"的话题。我们客观地来比。
5.1 核心设计理念的差异
| 维度 | Spark | Flink |
|---|---|---|
| 核心模型 | 批处理为核心,流是微批 | 流处理为核心,批是有界流 |
| 世界观 | 世界是由一批一批的数据组成的 | 世界是由一条一条的事件流组成的 |
| 延迟 | 秒级(微批) | 毫秒级(真流) |
| 状态 | 流处理状态较弱(Structured Streaming 在改进) | 状态管理是核心强项 |
| 批处理 | 出身就是批处理,非常成熟 | 批流一体,批处理也很强 |
💡 本质区别: Spark 是"批处理引擎,扩展支持流处理" Flink 是"流处理引擎,扩展支持批处理"
这决定了它们的基因和优势领域。
5.2 详细对比
性能对比
| 指标 | Spark | Flink |
|---|---|---|
| 批处理吞吐量 | 很高 | 很高(两者差不多,Spark 可能略高) |
| 流处理吞吐量 | 高 | 很高 |
| 延迟 | 秒级(Structured Streaming 可到 100ms) | 毫秒级(可到 10ms 以下) |
| 资源利用率 | 高 | 高 |
API 与生态
| 维度 | Spark | Flink |
|---|---|---|
| 语言支持 | Scala/Java/Python/R/SQL | Scala/Java/Python/SQL |
| SQL 能力 | 强(Catalyst 优化器成熟) | 强(持续改进中) |
| 机器学习 | MLlib(成熟) | Flink ML(较弱) |
| 图计算 | GraphX(维护中) | Gelly(维护中) |
| 生态系统 | 非常成熟,组件多 | 快速发展,生态相对小 |
| 社区活跃度 | 非常高 | 高(增长快) |
| 国内使用 | 非常广泛(几乎所有大数据公司) | 越来越多(互联网公司用得多) |
流处理能力(Flink 的强项)
| 能力 | Spark Structured Streaming | Flink |
|---|---|---|
| 事件时间 | 支持(Watermark) | 原生支持,非常完善 |
| 窗口 | 滚动、滑动、会话 | 滚动、滑动、会话、自定义窗口 |
| 状态管理 | 有(RocksDB 支持) | 非常强大(State Backend 成熟) |
| 状态后端 | HDFS / RocksDB | Memory / Fs / RocksDB |
| Exactly-Once | 支持(部分 Sink) | 原生支持(两阶段提交) |
| 迟到数据处理 | Watermark | Watermark + 侧输出 |
| CEP(复杂事件处理) | 不支持(需自己实现) | 原生支持(Flink CEP) |
| 背压 | 有(基于微批天然背压) | 有(基于 Credit 协议) |
批处理能力(Spark 的强项)
| 能力 | Spark | Flink |
|---|---|---|
| 批处理 API | DataFrame / SQL / RDD | DataStream / Table API / SQL |
| 优化器 | Catalyst(非常成熟) | 优化器在改进中 |
| 向量化执行 | 有(Spark 3.x 改进中) | 有(持续改进) |
| 机器学习 | MLlib(成熟) | Flink ML(较弱) |
| 图计算 | GraphX / GraphFrames | Gelly |
| 生态集成 | 非常丰富 | 丰富 |
5.3 怎么选?—— 场景驱动的选型
选 Spark 的场景
- ✅ 主要是批处理(ETL、数据仓库、离线分析)
- ✅ 需要机器学习(MLlib)
- ✅ 需要交互式查询(Spark SQL / Spark Thrift Server)
- ✅ 公司已经有 Hadoop/Spark 生态
- ✅ 流处理需求不高(分钟级延迟就够)
- ✅ 团队熟悉 Spark,招聘容易
选 Flink 的场景
- ✅ 主要是流处理(实时数仓、实时推荐、实时监控)
- ✅ 需要毫秒级延迟
- ✅ 需要复杂的事件时间处理
- ✅ 需要强大的状态管理(大状态、复杂状态)
- ✅ 需要 CEP(复杂事件处理)
- ✅ 需要 exactly-once 语义保证
- ✅ 公司是互联网/实时业务驱动
两者都用的场景(很常见)
- 批处理用 Spark,流处理用 Flink
- 数据湖 + Spark 做离线 ETL + Flink 做实时计算
- 这是目前很多中大型公司的实际架构
🎯 我的建议(给有经验的工程师):
- 两个都学,这是大数据工程师的标配
- 先学 Spark(生态好、工作多、入门容易、批处理是基础)
- 再学 Flink(流处理更强,互联网公司刚需)
- 不要站队,技术是工具,场景决定选型
- 批处理选 Spark,流处理选 Flink,这是目前的行业共识(但两者都在互相渗透)
5.4 未来趋势:批流一体
两个框架都在往"批流一体"方向走:
- Spark:Structured Streaming 越来越强,流处理能力在追赶
- Flink:批处理能力在完善,Table API/SQL 在统一批和流
未来的趋势是:一套 API,既能跑批也能跑流,用户不用关心底层是批还是流。 谁能先做到真正的批流一体,谁就能占据更大的市场。
六、Spark vs MapReduce / Storm / Tez:代际对比
6.1 三代计算引擎的演进
| 代际 | 代表 | 核心思想 | 延迟 | 吞吐量 |
|---|---|---|---|---|
| 第一代 | MapReduce | 两阶段批处理,磁盘落盘 | 分钟~小时 | 中 |
| 第二代 | Tez / Spark | DAG 执行,内存计算 | 秒~分钟 | 高 |
| 第三代 | Flink | 真流处理,事件时间 | 毫秒 | 很高 |
6.2 Spark vs MapReduce
| 对比项 | MapReduce | Spark |
|---|---|---|
| 计算模型 | Map + Reduce 两阶段 | DAG 多阶段 |
| 中间结果 | 写 HDFS | 放内存(可溢写磁盘) |
| 速度 | 慢(基准) | 快 10~100 倍 |
| API | Java 为主,代码量大 | Scala/Python/Java/SQL,简洁 |
| 流处理 | 不支持 | 支持(微批) |
| SQL | Hive(底层 MR) | Spark SQL(原生) |
| 机器学习 | Mahout(不活跃) | MLlib(活跃) |
| 内存需求 | 低 | 高 |
| 适用场景 | 简单的大批量处理 | 几乎所有大数据场景 |
🎯 MapReduce 还有必要学吗? 作为历史了解一下就行,新项目不会用了。 但理解 MapReduce 的思想(移动计算不移动数据、分而治之)对理解分布式计算有帮助。
6.3 Spark vs Storm
| 对比项 | Storm | Spark Streaming | Flink |
|---|---|---|---|
| 模型 | 逐条处理 | 微批 | 真流 |
| 延迟 | 毫秒级 | 秒级 | 毫秒级 |
| 吞吐量 | 较低 | 高 | 很高 |
| 状态管理 | 弱(需自己实现) | 有 | 强 |
| 容错 | Ack 机制 | 血缘 | Checkpoint |
| Exactly-Once | 不支持(Trident 支持但慢) | 支持 | 原生支持 |
| API | 较底层 | 高层 | 高层 |
| 现状 | 逐渐被替代 | 主流 | 越来越主流 |
🎯 Storm 还有必要学吗? 基本不用了。流处理要么用 Spark Streaming(简单场景),要么用 Flink(复杂场景)。 Storm 是"上一代"流计算框架,新项目不会选。
6.4 Spark vs Tez
Tez 是 Hadoop 生态里的 DAG 计算引擎,主要作为 Hive/Pig 的底层执行引擎。
| 对比项 | Tez | Spark |
|---|---|---|
| 定位 | Hive/Pig 的执行引擎 | 通用计算引擎 |
| API | 没有用户 API,只做执行引擎 | 有完整的用户 API |
| 生态 | 绑定 Hadoop 生态 | 独立,支持多种集群管理器 |
| 内存计算 | 有(比 MR 好) | 更成熟 |
| 流处理 | 不支持 | 支持 |
| 使用方式 | 一般通过 Hive 间接使用 | 直接写 Spark 程序 |
💡 Tez 的现状: Tez 主要用在 Hive on Tez,作为 Hive 的执行引擎(比 Hive on MR 快)。 但越来越多的 Hive 集群开始用 Spark 或 Tez 作为执行引擎,两者并存。 如果你不是专门做 Hive 运维,不用深入研究 Tez。
七、Spark SQL vs Hive vs Presto/Trino vs ClickHouse:SQL 引擎选型
大数据领域的 SQL 引擎很多,经常搞混。这里做一个全面对比。
7.1 各引擎定位
| 引擎 | 类型 | 核心定位 | 延迟 | 并发 |
|---|---|---|---|---|
| Spark SQL | 批处理/交互式 | 通用计算,ETL + 查询 | 秒~分钟 | 中 |
| Hive | 批处理 | 数据仓库,SQL 接口 | 分钟~小时 | 高 |
| Presto/Trino | 交互式查询 | 即席查询,联邦查询 | 秒级 | 中高 |
| ClickHouse | OLAP 数据库 | 列式存储,极速分析 | 毫秒~秒 | 中 |
| Impala | MPP 查询 | Hadoop 上的 MPP 查询 | 秒级 | 中 |
| Doris | OLAP 数据库 | 实时分析,MPP 架构 | 毫秒~秒 | 中高 |
| Druid | 时序 OLAP | 实时时序数据分析 | 毫秒~秒 | 中 |
7.2 详细对比
Spark SQL
- 优点:通用,既能 ETL 又能查询;生态好;和 Spark 其他组件集成
- 缺点:延迟较高(秒级起);不适合高并发点查
- 适用:ETL、批量分析、中等并发的即席查询
Hive
- 优点:数据仓库标准;元数据管理完善;并发能力强;生态最成熟
- 缺点:慢(底层 MR/Tez/Spark);不适合交互式查询
- 适用:离线数据仓库、大批量 ETL、低并发高吞吐的场景
Presto / Trino
- 优点:查询快(秒级);支持联邦查询(跨数据源 join);ANSI SQL 兼容好
- 缺点:不适合长时间 ETL;内存密集型,大查询容易 OOM;没有存储
- 适用:即席查询、跨数据源查询、数据探索
💡 Presto vs Trino 的关系: Presto 原本是 Facebook 开源的,后来核心团队分叉出 Trino(因为和 Facebook 的治理分歧)。 现在两者都在发展,Trino 更活跃一些。功能上大同小异。
ClickHouse
- 优点:极快(列式存储 + 向量化执行 + MPP);支持高吞吐写入;SQL 兼容好
- 缺点:不支持事务;join 能力弱;不适合频繁更新;运维有一定门槛
- 适用:OLAP 分析、日志分析、用户行为分析、实时报表
Doris
- 优点:MPP 架构,查询快;支持实时写入;兼容 MySQL 协议;运维相对简单
- 缺点:生态不如 Spark/Hive 成熟;超大规模集群案例相对少
- 适用:实时数仓、交互式分析、报表系统
7.3 选型决策
你的需求是什么?
│
├── 离线 ETL / 数据仓库建设
│ ├── 已经有 Hadoop 生态 → Hive(元数据) + Spark/Tez(执行)
│ └── 新建 → Spark SQL(直接做 ETL 和查询)
│
├── 即席查询 / 数据探索
│ ├── 数据在 Hadoop 上 → Presto/Trino 或 Spark SQL
│ ├── 需要跨数据源 → Presto/Trino(联邦查询)
│ └── 追求极致速度 → ClickHouse / Doris
│
├── 实时分析 / 实时报表
│ └── ClickHouse / Doris / Druid
│
├── 高并发点查(毫秒级,按 key 查)
│ └── 不是这些的强项,用 HBase / Redis
│
└── 机器学习 / 复杂数据处理
└── Spark(MLlib + DataFrame)🎯 实际生产中的常见组合:
- 离线数仓:Hive(存储+元数据) + Spark/Tez(执行)
- 实时数仓:Kafka + Flink + Doris/ClickHouse
- 即席查询:Presto/Trino 查 Hive 表
- 数据科学:Spark + Jupyter
没有一个引擎能搞定所有场景,都是组合使用。
八、生态集成:Spark 与 Hive/Kafka/HBase/Delta Lake 的协作
Spark 很少单独使用,通常是和其他大数据组件配合。这里讲最常见的几个集成。
8.1 Spark + Hive:最经典的组合
为什么要集成?
- Hive 有完善的元数据管理(库、表、分区、分桶)
- Hive 表是数据仓库的标准存储格式
- Spark 可以直接读 Hive 表,用 Spark SQL 查,比 Hive 本身快
怎么集成?
- 把
hive-site.xml放到 Spark 的conf目录 - 启动 Spark Session 时开启 Hive 支持:
val spark = SparkSession.builder()
.appName("Example")
.enableHiveSupport() // 开启 Hive 支持
.getOrCreate()- 直接用 SQL 查 Hive 表:
spark.sql("SELECT * FROM mydb.mytable WHERE dt = '2026-08-11'").show()Spark 写 Hive 表
// 覆盖写
df.write.mode("overwrite").saveAsTable("mydb.mytable")
// 追加写
df.write.mode("append").saveAsTable("mydb.mytable")
// 写分区表
df.write.partitionBy("dt").mode("overwrite").saveAsTable("mydb.mytable")⚠️ 注意: Spark 写 Hive 表时,默认格式是 Parquet(列式存储,比 TextFile 快很多)。 如果要和 Hive 兼容,确保两边的格式一致。
8.2 Spark + Kafka:流处理的标准搭配
Structured Streaming 读 Kafka
val df = spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "broker1:9092,broker2:9092")
.option("subscribe", "my-topic")
.option("startingOffsets", "latest")
.load()
// 解析 value(默认是二进制)
val parsed = df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")写 Kafka
df.writeStream
.format("kafka")
.option("kafka.bootstrap.servers", "broker1:9092")
.option("topic", "output-topic")
.option("checkpointLocation", "/path/to/checkpoint")
.start()关键配置
| 参数 | 说明 |
|---|---|
subscribe | 订阅的 topic,支持逗号分隔多个 |
subscribePattern | 正则匹配 topic |
startingOffsets | 起始位置:earliest / latest / 具体 offset |
maxOffsetsPerTrigger | 每次触发最多读多少条(限速) |
checkpointLocation | 检查点路径(必须设置,用于容错) |
💡 Kafka + Spark 的常见架构: 数据源 → Kafka → Spark Structured Streaming → 处理 → 写入 Hive/ClickHouse/Redis
这是实时数仓的标准架构之一。
8.3 Spark + HBase:批量读写
读 HBase
Spark 读 HBase 有几种方式:
- TableInputFormat(老方式,RDD API)
- HBase-Spark Connector(Cloudera 出的,DataFrame API)
- Phoenix + Spark(通过 Phoenix JDBC 读)
写 HBase
批量写 HBase 时注意:
- 用
saveAsNewAPIHadoopDataset或 BulkLoad - BulkLoad 更快(直接生成 HFile,不走 WAL),适合大批量导入
- 注意 RowKey 设计,避免写热点
⚠️ Spark + HBase 的坑:
- HBase 的 RegionServer 数量有限,Spark 并发太高会把 HBase 压垮
- 写 HBase 要控制并发度,或者用 BulkLoad
- 读 HBase 时如果全表扫描,会很慢,尽量用 RowKey 范围过滤
8.4 Spark + Delta Lake / Iceberg / Hudi:湖仓一体
这是当前最火的方向——在数据湖上提供 ACID 事务、时间旅行、增量处理等能力。
三者对比
| 特性 | Delta Lake | Apache Iceberg | Apache Hudi |
|---|---|---|---|
| 出品方 | Databricks | Netflix(后捐 Apache) | Uber(后捐 Apache) |
| ACID 事务 | ✅ | ✅ | ✅ |
| Schema 演进 | ✅ | ✅ | ✅ |
| 时间旅行 | ✅ | ✅ | ✅ |
| 增量读取 | ✅ | ✅ | ✅ |
| Upsert(更新/插入) | ✅(MERGE INTO) | ✅(MERGE INTO) | ✅(原生支持,更强) |
| 流式写入 | ✅ | ✅ | ✅ |
| 流式读取 | ✅ | ✅ | ✅ |
| 计算引擎支持 | Spark 最好,其他在加 | 多引擎支持(Spark/Flink/Trino) | Spark/Flink 都好 |
| 社区活跃度 | 高 | 高 | 高 |
| 国内使用 | 越来越多 | 越来越多(腾讯/阿里用得多) | 较多(阿里/字节用得多) |
Spark 用 Delta Lake 示例
// 读 Delta 表
val df = spark.read.format("delta").load("/path/to/table")
// 写 Delta 表
df.write.format("delta").mode("overwrite").save("/path/to/table")
// MERGE INTO(Upsert)
deltaTable.as("t")
.merge(updates.as("s"), "t.id = s.id")
.whenMatched().updateAll()
.whenNotMatched().insertAll()
.execute()
// 时间旅行
val oldDF = spark.read.format("delta")
.option("versionAsOf", "2026-08-01")
.load("/path/to/table")🎯 湖仓一体的意义: 以前数据仓库和数据湖是两套系统,数据搬来搬去,又贵又麻烦。 湖仓一体用一套存储(数据湖)+ 一层表格式(Delta/Iceberg/Hudi), 既能存原始数据,又能提供数仓的 ACID、性能、管理能力。 Spark 是湖仓一体的核心计算引擎之一。
中篇小结
这篇讲了 Spark 的技术栈和生态对比:
- Spark SQL:不只是 SQL 接口,是 Spark 2.0 后的核心 API 层;Catalyst 优化器 + Tungsten 执行引擎 + AQE 是它快的秘密
- Structured Streaming:新一代流计算,基于 Spark SQL,支持事件时间、Watermark、有状态计算,替代老的 DStream
- MLlib:分布式传统机器学习,不是深度学习框架;Pipeline 是核心概念;和 TensorFlow/PyTorch 互补
- GraphX:图计算组件,维护不活跃,简单场景够用,复杂场景用专门的图数据库
- Spark vs Flink:批处理选 Spark,流处理选 Flink;两者都学,场景驱动选型
- 代际对比:MapReduce/Storm 是上一代,新项目不用;Tez 主要做 Hive 执行引擎
- SQL 引擎选型:没有万能引擎,ETL 用 Spark/Hive,即席查询用 Presto/Trino,实时 OLAP 用 ClickHouse/Doris
- 生态集成:和 Hive/Kafka/HBase/Delta Lake 的集成是生产环境的标配
下篇预告:生产环境部署与最佳实践、性能调优圣经、排错手册、前沿方向(Spark Connect、Photon、GPU、大模型时代的 Spark)、学习路径与职业发展。
文档版本:V1.0(中篇) 最后更新:2026 年 8 月 适用 Spark 版本:3.5.x