Skip to content

Spark 完全入门指南(中篇):技术栈与生态对比 ​

上篇讲了 Spark 的核心原理,这篇讲 Spark 的五大组件,以及和其他大数据技术的深度对比。

面向有经验的工程师,不只是"是什么",更讲"为什么这么设计"、"和其他技术比怎么样"、"什么场景该用什么"。


目录 ​


一、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:到底用哪个? ​

特性RDDDataFrameDataset
类型安全✅ 编译时检查❌ 运行时才发现✅ 编译时检查
优化器❌ 无✅ 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 是开源的、可扩展的,你可以自己加优化规则。

怎么看执行计划? ​

scala
// 看逻辑计划和物理计划
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 换成 BroadcastJoinJoin 性能大幅提升
动态优化数据倾斜运行时发现倾斜的分区,自动拆分解决数据倾斜导致的长尾 Task

怎么开启? ​

properties
# 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 的流计算方案,基于微批处理。它的问题:

  1. API 不统一:DStream API 和 DataFrame API 是两套,批处理和流处理代码不能复用
  2. 延迟较高:最小批处理间隔 0.5 秒,实际一般 1~5 秒
  3. 事件时间支持弱:只能用处理时间,不支持基于事件时间的窗口
  4. 状态管理难:有状态计算要自己实现,容错复杂
  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 定义了"允许数据延迟多久"。

scala
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
基于RDDDataFrame/Dataset
APIDStream 专用 API和批处理统一
延迟秒级(最小 0.5s)秒级(可到 100ms)
事件时间不支持原生支持 + Watermark
状态管理需自己实现内置支持
语义保证at-least-onceexactly-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.mllibRDD维护模式不推荐
spark.mlDataFrame活跃开发官方推荐

spark.ml 的核心概念是 Pipeline(工作流):

DataFrame → Transformer → Transformer → Estimator → Model → 预测
  • Transformer:把一个 DataFrame 变成另一个 DataFrame(比如特征转换)
  • Estimator:训练数据,产生一个 Transformer(比如算法训练出模型)
  • Pipeline:把多个 Transformer 和 Estimator 串起来
scala
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 的好处:

  1. 特征工程和模型训练串在一起,不会漏步骤
  2. 可以一起保存/加载,部署方便
  3. 可以做交叉验证和参数搜索(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 的局限与替代方案 ​

局限 ​

  1. 算法更新慢:很多新算法(XGBoost、LightGBM、CatBoost)没有原生支持
  2. 不支持深度学习:MLP 只是简单的前馈网络,没有 CNN/RNN/Transformer
  3. 超参搜索慢:网格搜索是串行的,没有高级的超参优化
  4. 模型可解释性工具少:没有 SHAP、LIME 等工具
  5. 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 现在比较尴尬:

  1. 维护不活跃:GraphX 很久没有大的功能更新了
  2. API 基于 RDD:没有迁移到 DataFrame,享受不到 Catalyst/Tungsten 的优化
  3. 算法少:只有基础算法,很多高级算法没有
  4. 性能一般:和专门的图计算引擎比有差距

4.3 替代方案 ​

方案类型特点适用场景
GraphFramesSpark 上的图库基于 DataFrame,API 更现代,算法更多已经在用 Spark,需要图计算
Neo4j图数据库最流行的图数据库,查询语言 Cypher 强大图数据为主,需要频繁查询
JanusGraph分布式图数据库基于 HBase/Cassandra,可扩展超大规模图数据
DGL / PyG图神经网络框架深度学习 + 图,做 GNN图神经网络、节点预测
Flink GellyFlink 的图组件和 GraphX 类似,也不太活跃已经在用 Flink

🎯 建议:

  • 如果只是偶尔做简单的图计算(比如 PageRank),用 GraphX 或 GraphFrames 就行
  • 如果图是核心业务,用专门的图数据库(Neo4j/JanusGraph)
  • 如果要做图神经网络,用 DGL/PyG
  • GraphX 了解基本概念就行,不用深入研究

五、Spark vs Flink:大数据领域的"世纪之争" ​

这是大数据领域最常被问到的对比,也是最容易产生"信仰之争"的话题。我们客观地来比。

5.1 核心设计理念的差异 ​

维度SparkFlink
核心模型批处理为核心,流是微批流处理为核心,批是有界流
世界观世界是由一批一批的数据组成的世界是由一条一条的事件流组成的
延迟秒级(微批)毫秒级(真流)
状态流处理状态较弱(Structured Streaming 在改进)状态管理是核心强项
批处理出身就是批处理,非常成熟批流一体,批处理也很强

💡 本质区别: Spark 是"批处理引擎,扩展支持流处理" Flink 是"流处理引擎,扩展支持批处理"

这决定了它们的基因和优势领域。

5.2 详细对比 ​

性能对比 ​

指标SparkFlink
批处理吞吐量很高很高(两者差不多,Spark 可能略高)
流处理吞吐量高很高
延迟秒级(Structured Streaming 可到 100ms)毫秒级(可到 10ms 以下)
资源利用率高高

API 与生态 ​

维度SparkFlink
语言支持Scala/Java/Python/R/SQLScala/Java/Python/SQL
SQL 能力强(Catalyst 优化器成熟)强(持续改进中)
机器学习MLlib(成熟)Flink ML(较弱)
图计算GraphX(维护中)Gelly(维护中)
生态系统非常成熟,组件多快速发展,生态相对小
社区活跃度非常高高(增长快)
国内使用非常广泛(几乎所有大数据公司)越来越多(互联网公司用得多)
能力Spark Structured StreamingFlink
事件时间支持(Watermark)原生支持,非常完善
窗口滚动、滑动、会话滚动、滑动、会话、自定义窗口
状态管理有(RocksDB 支持)非常强大(State Backend 成熟)
状态后端HDFS / RocksDBMemory / Fs / RocksDB
Exactly-Once支持(部分 Sink)原生支持(两阶段提交)
迟到数据处理WatermarkWatermark + 侧输出
CEP(复杂事件处理)不支持(需自己实现)原生支持(Flink CEP)
背压有(基于微批天然背压)有(基于 Credit 协议)

批处理能力(Spark 的强项) ​

能力SparkFlink
批处理 APIDataFrame / SQL / RDDDataStream / Table API / SQL
优化器Catalyst(非常成熟)优化器在改进中
向量化执行有(Spark 3.x 改进中)有(持续改进)
机器学习MLlib(成熟)Flink ML(较弱)
图计算GraphX / GraphFramesGelly
生态集成非常丰富丰富

5.3 怎么选?—— 场景驱动的选型 ​

选 Spark 的场景 ​

  • ✅ 主要是批处理(ETL、数据仓库、离线分析)
  • ✅ 需要机器学习(MLlib)
  • ✅ 需要交互式查询(Spark SQL / Spark Thrift Server)
  • ✅ 公司已经有 Hadoop/Spark 生态
  • ✅ 流处理需求不高(分钟级延迟就够)
  • ✅ 团队熟悉 Spark,招聘容易
  • ✅ 主要是流处理(实时数仓、实时推荐、实时监控)
  • ✅ 需要毫秒级延迟
  • ✅ 需要复杂的事件时间处理
  • ✅ 需要强大的状态管理(大状态、复杂状态)
  • ✅ 需要 CEP(复杂事件处理)
  • ✅ 需要 exactly-once 语义保证
  • ✅ 公司是互联网/实时业务驱动

两者都用的场景(很常见) ​

  • 批处理用 Spark,流处理用 Flink
  • 数据湖 + Spark 做离线 ETL + Flink 做实时计算
  • 这是目前很多中大型公司的实际架构

🎯 我的建议(给有经验的工程师):

  1. 两个都学,这是大数据工程师的标配
  2. 先学 Spark(生态好、工作多、入门容易、批处理是基础)
  3. 再学 Flink(流处理更强,互联网公司刚需)
  4. 不要站队,技术是工具,场景决定选型
  5. 批处理选 Spark,流处理选 Flink,这是目前的行业共识(但两者都在互相渗透)

5.4 未来趋势:批流一体 ​

两个框架都在往"批流一体"方向走:

  • Spark:Structured Streaming 越来越强,流处理能力在追赶
  • Flink:批处理能力在完善,Table API/SQL 在统一批和流

未来的趋势是:一套 API,既能跑批也能跑流,用户不用关心底层是批还是流。 谁能先做到真正的批流一体,谁就能占据更大的市场。


六、Spark vs MapReduce / Storm / Tez:代际对比 ​

6.1 三代计算引擎的演进 ​

代际代表核心思想延迟吞吐量
第一代MapReduce两阶段批处理,磁盘落盘分钟~小时中
第二代Tez / SparkDAG 执行,内存计算秒~分钟高
第三代Flink真流处理,事件时间毫秒很高

6.2 Spark vs MapReduce ​

对比项MapReduceSpark
计算模型Map + Reduce 两阶段DAG 多阶段
中间结果写 HDFS放内存(可溢写磁盘)
速度慢(基准)快 10~100 倍
APIJava 为主,代码量大Scala/Python/Java/SQL,简洁
流处理不支持支持(微批)
SQLHive(底层 MR)Spark SQL(原生)
机器学习Mahout(不活跃)MLlib(活跃)
内存需求低高
适用场景简单的大批量处理几乎所有大数据场景

🎯 MapReduce 还有必要学吗? 作为历史了解一下就行,新项目不会用了。 但理解 MapReduce 的思想(移动计算不移动数据、分而治之)对理解分布式计算有帮助。

6.3 Spark vs Storm ​

对比项StormSpark StreamingFlink
模型逐条处理微批真流
延迟毫秒级秒级毫秒级
吞吐量较低高很高
状态管理弱(需自己实现)有强
容错Ack 机制血缘Checkpoint
Exactly-Once不支持(Trident 支持但慢)支持原生支持
API较底层高层高层
现状逐渐被替代主流越来越主流

🎯 Storm 还有必要学吗? 基本不用了。流处理要么用 Spark Streaming(简单场景),要么用 Flink(复杂场景)。 Storm 是"上一代"流计算框架,新项目不会选。

6.4 Spark vs Tez ​

Tez 是 Hadoop 生态里的 DAG 计算引擎,主要作为 Hive/Pig 的底层执行引擎。

对比项TezSpark
定位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交互式查询即席查询,联邦查询秒级中高
ClickHouseOLAP 数据库列式存储,极速分析毫秒~秒中
ImpalaMPP 查询Hadoop 上的 MPP 查询秒级中
DorisOLAP 数据库实时分析,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 本身快

怎么集成? ​

  1. 把 hive-site.xml 放到 Spark 的 conf 目录
  2. 启动 Spark Session 时开启 Hive 支持:
scala
val spark = SparkSession.builder()
  .appName("Example")
  .enableHiveSupport()  // 开启 Hive 支持
  .getOrCreate()
  1. 直接用 SQL 查 Hive 表:
scala
spark.sql("SELECT * FROM mydb.mytable WHERE dt = '2026-08-11'").show()

Spark 写 Hive 表 ​

scala
// 覆盖写
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 ​

scala
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 ​

scala
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 有几种方式:

  1. TableInputFormat(老方式,RDD API)
  2. HBase-Spark Connector(Cloudera 出的,DataFrame API)
  3. 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 LakeApache IcebergApache Hudi
出品方DatabricksNetflix(后捐 Apache)Uber(后捐 Apache)
ACID 事务✅✅✅
Schema 演进✅✅✅
时间旅行✅✅✅
增量读取✅✅✅
Upsert(更新/插入)✅(MERGE INTO)✅(MERGE INTO)✅(原生支持,更强)
流式写入✅✅✅
流式读取✅✅✅
计算引擎支持Spark 最好,其他在加多引擎支持(Spark/Flink/Trino)Spark/Flink 都好
社区活跃度高高高
国内使用越来越多越来越多(腾讯/阿里用得多)较多(阿里/字节用得多)

Spark 用 Delta Lake 示例 ​

scala
// 读 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 的技术栈和生态对比:

  1. Spark SQL:不只是 SQL 接口,是 Spark 2.0 后的核心 API 层;Catalyst 优化器 + Tungsten 执行引擎 + AQE 是它快的秘密
  2. Structured Streaming:新一代流计算,基于 Spark SQL,支持事件时间、Watermark、有状态计算,替代老的 DStream
  3. MLlib:分布式传统机器学习,不是深度学习框架;Pipeline 是核心概念;和 TensorFlow/PyTorch 互补
  4. GraphX:图计算组件,维护不活跃,简单场景够用,复杂场景用专门的图数据库
  5. Spark vs Flink:批处理选 Spark,流处理选 Flink;两者都学,场景驱动选型
  6. 代际对比:MapReduce/Storm 是上一代,新项目不用;Tez 主要做 Hive 执行引擎
  7. SQL 引擎选型:没有万能引擎,ETL 用 Spark/Hive,即席查询用 Presto/Trino,实时 OLAP 用 ClickHouse/Doris
  8. 生态集成:和 Hive/Kafka/HBase/Delta Lake 的集成是生产环境的标配

下篇预告:生产环境部署与最佳实践、性能调优圣经、排错手册、前沿方向(Spark Connect、Photon、GPU、大模型时代的 Spark)、学习路径与职业发展。


文档版本:V1.0(中篇) 最后更新:2026 年 8 月 适用 Spark 版本:3.5.x

基于 Vite 强力驱动 | 纯静态轻量托管