Skip to content

阿里云 EMR Apache Spark 入门学习 ​

一、Spark 基础认知与核心原理 ​

1.1 什么是 Apache Spark ​

Apache Spark 是一款统一的分布式计算引擎,专为大规模数据处理而设计,最初由加州大学伯克利分校 AMPLab 开发,目前是 Apache 顶级项目,也是大数据领域的事实标准计算引擎。

官方定义:Spark is a multi-language engine for executing data engineering, data science, and machine learning on single-node machines or clusters. 来源:Apache Spark 官网

Spark 的核心优势:

  1. 极致速度:基于内存计算,比传统 MapReduce 快 10~100 倍,支持 DAG(有向无环图)执行引擎,减少磁盘 IO。
  2. 统一栈:一套引擎覆盖批处理、SQL 查询、流计算、机器学习、图计算五大场景。
  3. 多语言支持:原生支持 Scala、Java、Python、R、SQL 多种开发语言。
  4. 生态兼容:兼容 Hadoop、Hive、YARN、Kafka 等所有大数据组件,可直接运行在 Hadoop 集群上。

1.2 Spark vs MapReduce 核心差异 ​

维度MapReduceSpark
计算模型磁盘计算,中间结果落盘内存计算,中间结果存在内存
执行模型只有 Map + Reduce 两阶段DAG 任意阶段,支持复杂任务流
延迟高,适合离线大任务低,支持交互式、迭代式计算
编程接口底层 API,开发效率低高层抽象(RDD/DataFrame),开发效率高
场景简单离线批处理批处理、交互式分析、机器学习、流计算

1.3 Spark 五大核心组件 ​

Spark 采用「一个核心引擎 + 多个上层组件」的架构:

  1. Spark Core:核心引擎,提供 RDD 抽象、调度、内存管理,是所有上层组件的基础。
  2. Spark SQL:结构化数据处理引擎,支持 SQL/HiveQL 查询,DataFrame/Dataset 抽象,是目前最常用的组件。
  3. Spark Streaming:流计算引擎,处理实时数据流,支持 Kafka、Flume 等数据源。
  4. MLlib:分布式机器学习库,提供常用的分类、回归、聚类、推荐算法。
  5. GraphX:分布式图计算引擎,用于社交网络、路径计算等图场景。

1.4 官方学习资源指引 ​


二、Spark 运行架构与模式 ​

2.1 核心运行组件 ​

Spark 采用主从(Master-Worker)架构,运行在 YARN 上时核心角色如下:

角色作用运行位置
Driver(驱动器)运行 main 方法,创建 SparkContext,负责任务调度、DAG 生成、任务分配Client 模式:提交节点;Cluster 模式:集群节点
Executor(执行器)工作进程,负责执行具体任务(Task),存储缓存数据所有 Core 工作节点
ApplicationMaster向 YARN 申请资源,管理 Executor 生命周期集群节点,每个应用一个
Cluster Manager集群资源管理器,负责资源分配YARN ResourceManager

2.2 YARN 两种运行模式 ​

你的 EMR 集群使用 YARN 作为资源管理器,Spark 支持两种提交模式:

(1)Client 模式(客户端模式) ​

  • Driver 运行在提交命令的本地节点(Master 节点)
  • 适合交互式开发、调试,能直接看到输出和日志
  • 对应命令:spark-shell --master yarn --deploy-mode client(你之前使用的模式)

(2)Cluster 模式(集群模式) ​

  • Driver 运行在集群的某个 Worker 节点上
  • 适合生产环境正式任务,提交后客户端可以断开连接
  • 对应命令:spark-submit --master yarn --deploy-mode cluster

2.3 任务执行完整流程 ​

  1. 提交应用,Driver 启动,创建 SparkContext 上下文
  2. ApplicationMaster 向 YARN 申请 Executor 资源
  3. YARN 在 Core 节点上启动 Executor 进程
  4. Driver 将代码解析为 DAG,切分为 Stage 和 Task
  5. Task 分发到 Executor 上执行,结果返回 Driver
  6. 任务结束,释放资源

三、入门入口:Spark Shell 交互式开发 ​

3.1 EMR 环境启动 Spark Shell ​

你的 EMR 集群已预装 Spark 3.5.3,无需任何配置,SSH 登录 Master 节点后直接执行启动命令:

sh
# 启动交互式终端,YARN客户端模式,和生产环境标准一致
spark-shell --master yarn --deploy-mode client

启动成功标志:

  • 出现 Spark 标志性 ASCII 艺术图案
  • 显示版本:version 3.5.3-emr
  • 命令提示符变为 scala>
  • 自动创建两个上下文:
    • sc:SparkContext,RDD 编程入口
    • spark:SparkSession,Spark SQL 编程入口

退出命令::quit

3.2 第一个 Spark 程序:词频统计(WordCount) ​

词频统计是大数据的「Hello World」,通过它可以快速理解 Spark 编程思想。

步骤1:准备数据 ​

在 HDFS 上创建测试文件:

sh
# 退出spark-shell后执行,或新开一个SSH窗口
echo "hello spark hello hadoop hello emr
spark is fast hadoop is stable
emr is easy to use" > words.txt

# 上传到HDFS
hdfs dfs -put words.txt /test/

hdfs dfs -cat /test/words.txt

image-20260907013212359

步骤2:Spark 实现词频统计 ​

回到 spark-shell 中执行:

scala
// 1. 从HDFS读取文件,创建RDD
val lines = sc.textFile("/test/words.txt")

// 2. 拆分单词,扁平化处理
val words = lines.flatMap(line => line.split(" "))

// 3. 每个单词映射为(单词, 1)键值对
val wordOne = words.map(word => (word, 1))

// 4. 按单词聚合,累加计数
val wordCount = wordOne.reduceByKey(_ + _)

// 5. 触发执行,收集结果并打印
wordCount.collect().foreach(println)

image-20260907013306875

步骤3:结果与原理说明 ​

输出结果:

(hello,3)
(spark,2)
(hadoop,2)
(emr,2)
(is,3)
(fast,1)
(stable,1)
(easy,1)
(to,1)
(use,1)

核心原理:

  • textFile、flatMap、map、reduceByKey 都是转换算子,不会立即执行,只记录操作逻辑
  • collect() 是行动算子,触发真正的计算,这是 Spark「惰性求值」的核心特性
  • 惰性求值的好处:Spark 可以在执行前优化整个执行计划,提升性能

三、 PySpark 交互命令行 ​

在终端中直接输入:

bash
pyspark

稍等几秒,出现 >>> 符号即表示进入了 Spark 环境。

3. 粘贴代码运行计算 ​

将下面代码复制并粘贴到 >>> 后面,回车:

python
# 1. 读取 HDFS 上的文件
text_file = sc.textFile("hdfs:///test/words.txt")

# 2. 切分单词并计算数量
counts = text_file.flatMap(lambda line: line.split(" ")) \
                 .map(lambda word: (word, 1)) \
                 .reduceByKey(lambda a, b: a + b)

# 3. 打印计算结果
print(counts.collect())
  • 期望输出: [('hello', 2), ('spark', 1), ('emr', 1)]
  • 输入 quit() 并回车,退出 PySpark 环境。

image-20260907030730527

第三步:用 Spark 操作数据库(Spark SQL 实操) ​

Spark SQL 可以像写 SQL 语句一样操作 Hive 表。

1. 创建一个 Python 脚本文件 ​

在终端输入:

bash
cat << 'EOF' > /root/demo_sql.py
from pyspark.sql import SparkSession

# 1. 初始化支持 Hive 的 Spark 实例
spark = SparkSession.builder.appName("Demo").enableHiveSupport().getOrCreate()

# 2. 创建数据库和表
spark.sql("CREATE DATABASE IF NOT EXISTS my_db")
spark.sql("CREATE TABLE IF NOT EXISTS my_db.users (id INT, name STRING)")

# 3. 插入数据
spark.sql("INSERT INTO my_db.users VALUES (1, 'Alice'), (2, 'Bob')")

# 4. 查询数据并打印
print("=== 查询结果 ===")
spark.sql("SELECT * FROM my_db.users").show()

spark.stop()
EOF

2. 运行该脚本 ​

bash
python3 /root/demo_sql.py
  • 期望输出: 终端会显示一个格式化好的表格,包含 Alice 和 Bob 的数据。

第四步:提交正式生产任务(spark-submit) ​

在实际工作中,我们不会一直在命令行里手动敲代码,而是将写好的代码用 spark-submit 提交给集群运行。

运行以下命令,将刚才的脚本提交给 YARN 资源管理器调度:

bash
spark-submit --master yarn --deploy-mode client /root/demo_sql.py
  • 这个命令会把任务发送给您的从节点(core-1-1 和 core-1-2)协同完成计算。

image-20260907031110564

参数说明:

  • --master yarn:使用集群的 YARN 资源管理器,和交互式环境的运行架构一致
  • --deploy-mode client:客户端模式,日志直接输出到当前终端,适合调试查看输出
  • 末尾填写脚本的绝对路径即可

总结:您只需要掌握这三个核心命令 ​

  1. pyspark:用来进入交互环境,边写代码边看结果(适合开发调试)。

  2. python3 xxx.py:用来在单机上快速运行写好的 Spark 脚本。

  3. spark-submit --master yarn xxx.py:用来将脚本提交到 Hadoop/YARN 集群真正进行分布式计算(适合生产环境)。


四、核心抽象:RDD 编程详解 ​

4.1 RDD 核心概念 ​

RDD(Resilient Distributed Dataset,弹性分布式数据集)是 Spark 最基础的抽象,代表一个不可变、可分区、里面的元素可并行计算的集合。

五大核心特性(官方定义):

  1. 分区(Partition):数据分片,是并行计算的基本单位
  2. 计算函数:每个分区都有一个计算函数
  3. 依赖关系:RDD 之间的血缘关系,记录转换过程
  4. 分区器:键值对 RDD 可指定分区策略
  5. 优先位置:计算任务调度到数据所在的节点,移动计算不移动数据

关键特性:

  • 不可变性:RDD 一旦创建就不能修改,每次转换都会生成新的 RDD
  • 弹性:数据丢失可以通过血缘关系重新计算,自动容错
  • 分布式:数据分散存储在集群多个节点上

4.2 RDD 的三种创建方式 ​

scala
// 1. 从集合创建(测试常用)
val rdd1 = sc.parallelize(List(1,2,3,4,5))

// 2. 从外部存储创建(生产常用)
val rdd2 = sc.textFile("/test/words.txt") // HDFS文件
val rdd3 = sc.textFile("hdfs:///test/words.txt") // 完整HDFS路径

// 3. 从已有RDD转换生成
val rdd4 = rdd1.map(_ * 2)

4.3 常用转换算子(Transformation) ​

转换算子返回新的 RDD,惰性执行,不会触发计算。

基础转换 ​

scala
val rdd = sc.parallelize(List(1,2,3,4,5,6))

// map:对每个元素做映射,一对一转换
rdd.map(_ * 2).collect() // 输出 2,4,6,8,10,12

// filter:过滤符合条件的元素
rdd.filter(_ > 3).collect() // 输出 4,5,6

// flatMap:映射后扁平化,一对多
val lines = sc.parallelize(List("hello spark", "hello hadoop"))
lines.flatMap(_.split(" ")).collect() // 输出 hello,spark,hello,hadoop

// distinct:去重
val rdd2 = sc.parallelize(List(1,2,2,3,3,3))
rdd2.distinct().collect() // 输出 1,2,3

键值对 RDD 转换(Pair RDD) ​

由 (key, value) 元组组成的 RDD,是聚合计算的核心。

scala
val pairRdd = sc.parallelize(List(("a",1), ("b",2), ("a",3), ("b",4)))

// reduceByKey:按key聚合
pairRdd.reduceByKey(_ + _).collect() // 输出 (a,4), (b,6)

// groupByKey:按key分组
pairRdd.groupByKey().collect() // 输出 (a,CompactBuffer(1,3)), (b,CompactBuffer(2,4))

// sortByKey:按key排序
pairRdd.sortByKey().collect() // 按字母升序

// mapValues:只对value做映射
pairRdd.mapValues(_ * 10).collect() // 输出 (a,10), (b,20)...

4.4 常用行动算子(Action) ​

行动算子触发真正的计算,返回结果或写入外部存储。

scala
val rdd = sc.parallelize(List(1,2,3,4,5))

// collect:收集所有元素到Driver端
rdd.collect()

// count:统计元素数量
rdd.count() // 5

// first:取第一个元素
rdd.first() // 1

// take:取前n个元素
rdd.take(3) // 1,2,3

// reduce:全局聚合
rdd.reduce(_ + _) // 15

// foreach:遍历每个元素
rdd.foreach(println)

// saveAsTextFile:保存到HDFS
rdd.saveAsTextFile("/test/output")

4.5 RDD 持久化与缓存 ​

Spark RDD 默认是惰性求值,每次行动算子都会从头重新计算,对于重复使用的 RDD,可以缓存到内存中,大幅提升性能。

scala
val rdd = sc.textFile("/test/words.txt").flatMap(_.split(" "))

// 缓存到内存(最常用)
rdd.cache()

// 完整指定存储级别
import org.apache.spark.storage.StorageLevel
rdd.persist(StorageLevel.MEMORY_ONLY)

// 清除缓存
rdd.unpersist()

适用场景:迭代计算、多次查询的热点数据,是 Spark 性能优化的核心手段之一。


五、Spark SQL 入门(当前主流开发方式) ​

5.1 Spark SQL 概述 ​

Spark SQL 是 Spark 用于结构化数据处理的模块,提供了 DataFrame/Dataset 高级抽象,支持标准 SQL 查询,性能比 RDD 更高(有 Catalyst 优化器),也是目前企业开发的主流方式。

核心抽象:

  • DataFrame:带 Schema 的分布式数据集,类似关系型数据库的表,行由 Row 对象组成
  • Dataset:强类型的 DataFrame,支持编译时类型检查,Scala/Java 推荐使用
  • SparkSession:Spark SQL 的统一入口,spark-shell 中已自动创建为 spark 变量

5.2 DataFrame 基础操作 ​

(1)创建 DataFrame ​

scala
// 从HDFS读取文本文件创建
val df = spark.read.text("/test/words.txt")

// 读取JSON文件
val jsonDF = spark.read.json("/test/data.json")

// 从RDD转换
val rdd = sc.parallelize(List(("张三", 20), ("李四", 21), ("王五", 19)))
val personDF = rdd.toDF("name", "age")

(2)DataFrame 常用 DSL 操作 ​

scala
// 查看Schema
personDF.printSchema()

// 查看数据
personDF.show()

// 选择列
personDF.select("name", "age").show()

// 条件过滤
personDF.filter($"age" > 20).show()

// 分组聚合
personDF.groupBy("age").count().show()

// 排序
personDF.orderBy($"age".desc).show()

5.3 与 Hive 整合(EMR 开箱即用) ​

你的 EMR 集群已默认配置好 Hive 元数据集成,Spark SQL 可以直接查询 Hive 中的表,无需任何额外配置。

scala
// 查看Hive中的数据库
spark.sql("show databases").show()

// 使用test数据库
spark.sql("use test_db")

// 查看表
spark.sql("show tables").show()

// 执行SQL查询
val result = spark.sql("select * from student where age > 20")
result.show()

原理:Spark SQL 复用 Hive 的 MetaStore 元数据,计算使用 Spark 引擎,比原生 Hive 快数倍到数十倍。

5.4 常用 Spark SQL 语法 ​

Spark SQL 兼容标准 SQL 和 HiveQL,常用语法和关系型数据库一致:

scala
-- 建表
spark.sql("""
create table if not exists user_behavior (
    user_id string,
    item_id string,
    click_time string,
    pv int
)
row format delimited fields terminated by ','
""")

-- 插入数据
spark.sql("insert into user_behavior values ('u001','i001','2026-09-06',5)")

-- 分组统计
spark.sql("""
select user_id, sum(pv) as total_pv 
from user_behavior 
group by user_id 
order by total_pv desc
""").show()

六、Spark Web UI 监控与调试 ​

6.1 两个核心 UI 入口 ​

结合你 EMR 集群的端口配置:

UI 类型地址作用生命周期
任务运行 UIhttp://Master公网IP:4040查看当前正在运行的任务详情仅任务运行期间存在
历史服务 UIhttp://Master公网IP:18080查看所有历史任务记录常驻服务,集群运行就存在

6.2 UI 核心功能 ​

  1. Jobs 页面:查看所有作业,每个作业对应一个行动算子
  2. Stages 页面:查看阶段划分,DAG 可视化,了解任务执行流程
  3. Storage 页面:查看缓存的 RDD 存储情况、内存占用
  4. Executors 页面:查看所有 Executor 的资源使用、任务数、GC 情况
  5. Environment 页面:查看 Spark 配置参数、依赖 Jar 包

6.3 简单性能排查 ​

  • 任务慢:看 Stages 页面,找到耗时最长的 Stage,查看数据倾斜、GC 情况
  • 内存不足:看 Executors 页面,查看堆内存使用,是否有 OOM 报错
  • 任务失败:点击任务详情,查看错误日志,定位代码或数据问题

七、入门实战案例:访问日志分析 ​

7.1 需求与数据准备 ​

需求:模拟 Nginx 访问日志,统计访问量 Top3 的 IP、不同状态码的数量。

步骤1:创建测试数据

sh
# 在HDFS创建日志文件
cat > access.log << EOF
192.168.1.1 200 GET /index.html
192.168.1.2 200 GET /about.html
192.168.1.1 404 GET /notfound
192.168.1.3 200 GET /index.html
192.168.1.1 200 GET /list.html
192.168.1.2 500 GET /api
192.168.1.3 200 GET /about.html
192.168.1.1 200 GET /index.html
EOF

hdfs dfs -put access.log /test/

7.2 RDD 实现 ​

scala
// 读取日志
val logs = sc.textFile("/test/access.log")

// 解析每行,提取IP和状态码
val logInfo = logs.map(line => {
  val arr = line.split(" ")
  (arr(0), arr(1)) // (IP, 状态码)
})

// 需求1:统计每个IP的访问量,取Top3
val ipCount = logInfo.map(x => (x._1, 1))
  .reduceByKey(_ + _)
  .map(x => (x._2, x._1))
  .sortByKey(false)
  .take(3)

println("访问量Top3 IP:")
ipCount.foreach(println)

// 需求2:统计各状态码数量
val statusCount = logInfo.map(x => (x._2, 1))
  .reduceByKey(_ + _)

println("状态码统计:")
statusCount.collect().foreach(println)

image-20260907013845002

7.3 Spark SQL 实现 ​

scala
// 创建DataFrame
val logDF = logs.map(line => {
  val arr = line.split(" ")
  (arr(0), arr(1), arr(2), arr(3))
}).toDF("ip", "status", "method", "url")

// 注册临时视图
logDF.createOrReplaceTempView("access_log")

// SQL查询Top3 IP
spark.sql("""
select ip, count(*) as cnt 
from access_log 
group by ip 
order by cnt desc 
limit 3
""").show()

// SQL查询状态码统计
spark.sql("""
select status, count(*) as num 
from access_log 
group by status
""").show()

image-20260907013909122


八、学习路径与进阶方向 ​

8.1 入门阶段(当前) ​

  • 掌握 RDD 核心概念、常用算子
  • 掌握 Spark SQL 基础查询、DataFrame 操作
  • 理解 Spark 运行架构、惰性求值、DAG 执行原理
  • 能独立完成简单的离线数据统计任务

8.2 进阶阶段 ​

  1. 核心原理:DAG 划分、Stage 划分原理、宽窄依赖、共享变量
  2. 性能调优:资源配置、缓存策略、数据倾斜优化、Shuffle 调优
  3. 流计算:Spark Streaming / Structured Streaming 实时数据处理
  4. 生态整合:和 Kafka、HBase、Hive、S3 的整合使用

8.3 高级阶段 ​

  1. Spark 内核原理、源码解析
  2. 企业级数仓建设、湖仓一体
  3. 机器学习 MLlib 实战
  4. 生产环境运维、监控、故障排查

九、使用注意事项(EMR 环境) ​

  1. 资源释放:学习结束后退出 spark-shell,释放集群资源;集群不用时及时释放,避免持续计费。
  2. 数据存储:重要数据存放在 HDFS 或 OSS,不要存在节点本地磁盘,释放集群会丢失。
  3. 提交方式:交互式开发用 spark-shell,正式任务用 spark-submit 提交脚本。
  4. 版本兼容:EMR 5.21.0 对应 Spark 3.5.3,代码和官方 Spark 3.x 完全兼容。

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