spark简介

Spark 全面解析:MapReduce 慢在哪,Spark 怎么改的,它到底怎么跑

MapReduce 能处理大数据,但它有两个硬伤:

  • 太慢:每轮 Map 结果都要写磁盘,Reduce 再读磁盘,一个复杂任务来回读写好几次
  • 太笨:每步都得排序,就算不需要排序也排;每个任务起一个 JVM,启动就得几秒

Spark 做了一件事:把中间结果放内存,用线程替代进程。

  • 内存计算 → 磁盘 IO 大幅减少
  • 线程复用 → 任务启动从秒级降到毫秒级
  • DAG 调度 → 多阶段任务不用重复读写

Spark 比 MapReduce 快在哪?

对比维度 MapReduce Spark
中间结果存储 磁盘(HDFS) 内存(优先),磁盘兜底
任务启动开销 每个任务起 JVM(秒级) 线程池复用(毫秒级)
排序策略 强制排序(不管需不需要) 按需排序(Hash 聚合不排)
多阶段任务 每阶段写一次磁盘 DAG 串联,内存流转
迭代计算 每次都重读 HDFS 内存缓存,迭代快速

典型场景: 机器学习训练要迭代 100 次,MapReduce 每次读磁盘,Spark 第一次读进去后面全在内存——速度差 10-100 倍。

核心概念:RDD 和 DStream

RDD(弹性分布式数据集):Spark 里一切数据的抽象。

  • 分布式:数据分片存在多台机器上
  • 不可变:创建后不能改,要改就生成新的 RDD
  • 惰性计算mapfilter 只是记下了”要做什么”,不会立刻算,等你要结果(collectcount)才真正执行
  • 容错:记录了”血缘关系”,某分区丢了可以从源头重算

DStream(离散流):Spark Streaming 里的数据流抽象。

本质是”按时间切成小块的 RDD 序列”——每 1 秒的微批,每个微批就是一个 RDD。

实时数据流 → 按时间切片 → 每片变成 RDD → 用 RDD 算子处理 → 输出

架构:谁负责什么

Spark 跑一个任务,涉及四个角色:

角色 做什么 在哪跑
Driver 跑你的 main 方法,解析任务,生成执行计划 客户端或集群里
Executor 真正干活的进程,执行具体的 Task Worker 节点上
Cluster Manager 管资源的(CPU、内存),分配 Worker 和 Executor Standalone/YARN/K8s
Worker Node 物理机器,运行 Executor 集群中的节点

一个任务的完整流程:

  1. 你提交代码 → Cluster Manager 分配资源
  2. 启动 Driver → Driver 里创建 SparkContext(上下文)
  3. Driver 申请 Executor → Cluster Manager 在 Worker 上启动 Executor
  4. Driver 把代码拆成 Task → 发给 Executor 执行
  5. Executor 执行完 → 结果返回 Driver

核心调度机制:DAG + Stage

Spark 的调度分两层:

1. DAG Scheduler(高层调度):

把你的代码(RDD 的转换链)画成一个有向无环图,然后按 Shuffle 边界切成 Stage。

  • 窄依赖:一个父 RDD 的分区只被一个子 RDD 分区用 → 可以合并在同一个 Stage,不用 Shuffle
  • 宽依赖:一个父 RDD 分区被多个子 RDD 分区用 → 必须切 Stage,中间要 Shuffle
RDD 转换链:
textFile → flatMap → map → reduceByKey → mapsave

切 Stage:
Stage 1(窄依赖):textFile → flatMap → map
Stage 2(宽依赖,reduceByKey 触发 Shuffle):reduceByKey → mapsave

2. Task Scheduler(底层调度):

把 Stage 拆成 Task(每个分区一个 Task),发给 Executor 去跑。哪个 Executor 有空就发给谁,失败了就重试。

运行模式:按环境选

模式 说明 什么时候用
Local 单机跑,在本地 JVM 里 开发调试、写单元测试
Standalone Spark 自带的集群管理 小集群、快速部署
YARN 跑在 Hadoop YARN 上 已有 Hadoop 集群,统一资源管理
Kubernetes 跑在 K8s 上 云原生环境,容器化部署
Mesos 跑在 Mesos 上 多框架共享集群(用得少了)

提交命令统一,只是 Master 参数不同:

# Local 模式
spark-submit --master local[4] myapp.jar

# Standalone 模式
spark-submit --master spark://master:7077 myapp.jar

# YARN 模式
spark-submit --master yarn --deploy-mode cluster myapp.jar

# K8s 模式
spark-submit --master k8s://https://k8s:6443 myapp.jar