一文讲透 Apache Spark:从核心概念、运行架构到性能优化
1. Spark 是什么
Apache Spark 是一个面向大规模数据处理的统一计算引擎。它可以处理离线批任务、交互式 SQL、流式计算、机器学习和图计算等场景。和传统 MapReduce 相比,Spark 更强调内存计算、DAG 执行计划和统一 API,因此在迭代计算、复杂 ETL、数据分析和近实时处理场景中更常见。
简单理解:Spark 不是一个存储系统,而是一个计算引擎。它通常和 HDFS、S3、Hive、Kafka、Iceberg、Delta Lake、HBase、MySQL、PostgreSQL 等系统配合使用,从外部系统读取数据,经过分布式计算后再写回目标存储。
版本说明:本文概念主要适用于 Spark 3.5.x 和 Spark 4.x。根据 Apache Spark 官方文档,截至 2026-08-16,官方 latest 文档指向 Spark 4.2.0。
2. Spark 解决了什么问题
在单机程序里,数据量小时可以直接用 SQL、Python、Java Stream 或脚本处理。但当数据规模达到几十 GB、TB 甚至 PB 级别时,单机 CPU、内存、磁盘 IO 和网络能力都会成为瓶颈。Spark 的核心价值是把一个大任务拆成很多小任务,分发到集群中的多台机器并行执行。
典型使用场景包括:
| 场景 | 说明 |
|---|---|
| 离线 ETL | 清洗日志、订单、用户行为、埋点数据 |
| 数仓加工 | ODS、DWD、DWS、ADS 分层计算 |
| 交互式分析 | 通过 Spark SQL 查询大规模数据 |
| 实时/准实时计算 | 消费 Kafka 数据并聚合、清洗、落库 |
| 机器学习 | 使用 MLlib 做特征处理、训练和预测 |
| 图计算 | 使用 GraphX 处理关系网络、传播路径等问题 |
3. Spark 的核心架构
一个 Spark 应用通常由 Driver、Cluster Manager 和 Executor 组成。
| 组件 | 作用 |
|---|---|
| Driver | 运行用户主程序,创建 SparkSession/SparkContext,生成执行计划,调度任务 |
| Cluster Manager | 负责资源管理,例如 Standalone、YARN、Kubernetes |
| Worker Node | 集群中的工作节点,用来运行 Executor |
| Executor | 真正执行计算任务的进程,负责运行 Task、缓存数据、返回结果 |
| Task | Spark 最小执行单元,通常对应一个分区上的计算 |
| Job | 由一个 Action 触发的完整计算任务 |
| Stage | Job 被 Shuffle 边界切分后的执行阶段 |
| DAG | Spark 根据转换操作生成的有向无环图执行计划 |
执行流程可以概括为:用户代码提交到 Driver,Driver 构建 DAG,DAG Scheduler 切分 Stage,Task Scheduler 把 Task 分发给 Executor,Executor 执行计算并把结果或状态返回给 Driver。
4. Spark 的几个核心抽象
4.1 RDD
RDD,全称 Resilient Distributed Dataset,是 Spark 最早的核心抽象。它表示一个可以分区、并行处理、具备容错能力的数据集合。RDD 支持两类操作:Transformation 和 Action。
Transformation 会生成新的 RDD,例如 map、filter、flatMap、reduceByKey。Action 会真正触发计算,例如 count、collect、first、saveAsTextFile。
RDD 的特点是控制力强,但需要开发者自己关心较多底层细节。今天的新项目通常更推荐优先使用 DataFrame 或 Spark SQL,只有在需要非常底层的控制时再使用 RDD。
4.2 DataFrame
DataFrame 是带 Schema 的分布式数据表,可以理解为 Spark 里的“分布式表格”。它有列名和列类型,适合结构化数据处理。相比 RDD,DataFrame 能让 Spark 获取更多结构信息,从而通过 Catalyst Optimizer 做执行计划优化。
例如:
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, sum as sum_
spark = SparkSession.builder.appName("order-summary").getOrCreate()
orders = spark.read.option("header", True).csv("/data/orders.csv")
result = (
orders
.filter(col("status") == "PAID")
.groupBy("user_id")
.agg(sum_("amount").alias("total_amount"))
)
result.write.mode("overwrite").parquet("/warehouse/order_summary")
4.3 Dataset
Dataset 主要出现在 Scala 和 Java API 中,它结合了 RDD 的类型安全和 Spark SQL 的优化能力。Scala 中的 DataFrame 本质上是 Dataset[Row]。Python 因为语言动态特性,没有和 Scala/Java 完全一样的 Dataset API。
4.4 Spark SQL
Spark SQL 允许你用 SQL 处理分布式数据。它既可以查询临时视图,也可以和 Hive Metastore、数据湖表格式、JDBC 数据源结合使用。
orders.createOrReplaceTempView("orders")
spark.sql("""
SELECT user_id, SUM(amount) AS total_amount
FROM orders
WHERE status = 'PAID'
GROUP BY user_id
""").show()
4.5 Structured Streaming
Structured Streaming 是 Spark 推荐的结构化流处理模型。它把不断到来的数据流抽象成一张持续追加的表,开发者可以用和批处理接近的 DataFrame/SQL 写法表达流式计算。
events = (
spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "localhost:9092")
.option("subscribe", "order_events")
.load()
)
query = (
events.writeStream
.format("console")
.option("checkpointLocation", "/tmp/spark-checkpoints/order-events")
.start()
)
query.awaitTermination()
需要注意:Structured Streaming 的端到端语义不仅取决于 Spark checkpoint,也取决于数据源和数据写入端是否支持幂等、事务或可恢复提交。
5. Spark 为什么快
Spark 快的原因不是单一的“内存计算”,而是一组机制共同作用:
| 机制 | 说明 |
|---|---|
| 内存缓存 | 对重复使用的数据可以 cache/persist,减少重复 IO |
| DAG 执行 | 多个转换可以合并成执行计划,减少不必要的落盘 |
| 延迟执行 | Transformation 不立即计算,只有 Action 才触发任务 |
| Catalyst 优化器 | 针对 DataFrame/Spark SQL 优化逻辑计划和物理计划 |
| 列式存储适配 | 适合读取 Parquet、ORC 等列式格式,只扫描需要的列 |
| AQE | Adaptive Query Execution 可在运行时优化 Shuffle 分区、Join 策略和数据倾斜 |
6. Spark 的执行模型
Spark 的 Transformation 是惰性执行的。例如:
df2 = df.filter(col("amount") > 100)
df3 = df2.groupBy("user_id").count()
上面两行不会立刻触发集群计算,只是描述了计算逻辑。直到调用:
df3.show()
Spark 才会真正生成 Job,把它拆成 Stage 和 Task 执行。
Shuffle 是 Spark 性能优化中最重要的概念之一。当数据需要按 key 重新分布时,比如 groupBy、join、distinct、repartition,就可能触发 Shuffle。Shuffle 会涉及网络传输、磁盘读写和排序,是最容易变慢、OOM 或数据倾斜的地方。
7. 常见部署模式
Spark 可以运行在本地,也可以运行在集群上。
| 模式 | 说明 |
|---|---|
| Local | 本地开发和调试 |
| Standalone | Spark 自带的简单集群管理模式 |
| YARN | Hadoop 生态中常见的资源管理方式 |
| Kubernetes | 云原生环境中常见的部署方式 |
还要区分 client mode 和 cluster mode。client mode 中 Driver 运行在提交任务的机器上,cluster mode 中 Driver 运行在集群内部。生产环境更常见的是 cluster mode,因为它更适合长任务和统一资源管理。
8. 性能优化建议
第一,优先使用 DataFrame 和 Spark SQL。它们能让 Spark 理解 Schema 和表达式,从而进行更多优化。除非确实需要底层控制,不要一开始就写 RDD。
第二,尽量减少 Shuffle。能提前过滤就提前过滤,能只选择必要字段就不要 select *,小表 Join 大表时可以考虑 broadcast join。
第三,合理控制分区。分区太少会导致并行度不足,分区太多会导致调度开销变大。常见参数是 spark.sql.shuffle.partitions,默认值不一定适合所有任务,需要结合数据量和集群资源调整。
第四,谨慎使用 cache()。只有中间结果会被多次复用时才值得缓存。缓存后不用了要 unpersist(),否则可能挤占执行内存。
第五,避免大数据量 collect()。collect() 会把分布式数据拉回 Driver,大数据量时容易导致 Driver OOM。调试时可以用 show()、take() 或采样。
第六,关注数据倾斜。如果少数 key 的数据量特别大,某些 Task 会明显慢于其他 Task。可以通过 AQE、加盐、拆分热点 key、预聚合等方式处理。
第七,使用合适的数据格式。分析型场景优先考虑 Parquet 或 ORC,通常比 CSV/JSON 更适合大规模查询,因为它们支持列式存储、压缩和谓词下推。
第八,认真看 Spark UI。默认情况下 Driver 会暴露 4040 Web UI,可以查看 Job、Stage、Task、SQL 执行计划、Shuffle、缓存、Executor 内存和失败原因。
9. 常见问题和坑
| 问题 | 原因 | 建议 |
|---|---|---|
| Driver OOM | collect() 拉回过多数据 | 改用分布式写出、分页查看或采样 |
| Executor OOM | 单分区数据过大、Join/聚合过重 | 调整分区、处理倾斜、优化 Join |
| 任务很慢 | Shuffle 大、数据倾斜、文件太碎 | 看 Spark UI,定位慢 Stage 和慢 Task |
| 小文件过多 | 上游写入分区过细 | 合并小文件,合理控制输出分区 |
| UDF 性能差 | 优化器难以理解普通 UDF 内部逻辑 | 优先使用内置函数,必要时再用 UDF |
| Streaming 重复消费 | checkpoint 不稳定或 sink 不幂等 | 固定 checkpoint 路径,写入端设计幂等 |
| 本地能跑集群失败 | 闭包、依赖、环境变量不同 | 用 spark-submit 管理依赖并检查集群环境 |
10. Spark 和其他技术的关系
Spark 和 Hadoop 不是替代关系。Hadoop 更常被理解为 HDFS、YARN、MapReduce 等组件的集合,而 Spark 可以运行在 YARN 上,也可以读写 HDFS。
Spark 和 Flink 都能做大数据计算。一般来说,Spark 在批处理、SQL、数仓 ETL、生态兼容方面非常成熟;Flink 在低延迟流处理、复杂事件处理和状态管理方面更有优势。实际选型要看团队经验、延迟要求、任务类型和现有平台。
Spark 和 Hive 也不是简单替代关系。Hive 更偏数仓元数据、SQL 生态和表管理,Spark SQL 可以查询 Hive 表,也可以作为 Hive SQL 之外的高性能计算引擎。
11. 推荐学习路线
- 先理解 Driver、Executor、Job、Stage、Task、Partition、Shuffle。
- 用 PySpark 或 Scala 写几个 DataFrame 批处理任务。
- 学会读 Spark UI,能定位慢 Stage、Shuffle 大小和失败 Task。
- 学 Spark SQL、Parquet、分区表、Join 优化和 AQE。
- 再学习 Structured Streaming、checkpoint、水位线、窗口聚合和幂等写入。
- 最后再深入内存管理、资源参数、数据倾斜治理和生产调优。
12. 总结
Spark 的本质是一个统一的大规模分布式计算引擎。它通过 Driver 统一调度,通过 Executor 并行执行,通过 DAG 和延迟执行优化计算过程,通过 DataFrame/Spark SQL 提供高级抽象和优化能力。
学习 Spark 不要只背 API,更要理解执行模型。真正影响生产任务稳定性的,往往是分区、Shuffle、Join、数据倾斜、缓存、文件格式、资源配置和外部系统写入语义。只要掌握这些核心概念,再结合 Spark UI 做定位,Spark 就不再是一个黑盒,而是一套可以分析、优化和治理的大数据计算体系。
参考资料
- Apache Spark Documentation: https://spark.apache.org/documentation
- Apache Spark 4.2.0 Overview: https://spark.apache.org/docs/latest/
- Spark Cluster Mode Overview: https://spark.apache.org/docs/latest/cluster-overview.html
- Spark RDD Programming Guide: https://spark.apache.org/docs/latest/rdd-programming-guide
- Spark SQL, DataFrames and Datasets Guide: https://spark.apache.org/docs/latest/sql-programming-guide
- Structured Streaming Programming Guide: https://spark.apache.org/docs/latest/streaming/index.html
- Spark SQL Performance Tuning: https://spark.apache.org/docs/latest/sql-performance-tuning
更多推荐
所有评论(0)