阿里云 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 的核心优势:
- 极致速度:基于内存计算,比传统 MapReduce 快 10~100 倍,支持 DAG(有向无环图)执行引擎,减少磁盘 IO。
- 统一栈:一套引擎覆盖批处理、SQL 查询、流计算、机器学习、图计算五大场景。
- 多语言支持:原生支持 Scala、Java、Python、R、SQL 多种开发语言。
- 生态兼容:兼容 Hadoop、Hive、YARN、Kafka 等所有大数据组件,可直接运行在 Hadoop 集群上。
1.2 Spark vs MapReduce 核心差异
| 维度 | MapReduce | Spark |
|---|---|---|
| 计算模型 | 磁盘计算,中间结果落盘 | 内存计算,中间结果存在内存 |
| 执行模型 | 只有 Map + Reduce 两阶段 | DAG 任意阶段,支持复杂任务流 |
| 延迟 | 高,适合离线大任务 | 低,支持交互式、迭代式计算 |
| 编程接口 | 底层 API,开发效率低 | 高层抽象(RDD/DataFrame),开发效率高 |
| 场景 | 简单离线批处理 | 批处理、交互式分析、机器学习、流计算 |
1.3 Spark 五大核心组件
Spark 采用「一个核心引擎 + 多个上层组件」的架构:
- Spark Core:核心引擎,提供 RDD 抽象、调度、内存管理,是所有上层组件的基础。
- Spark SQL:结构化数据处理引擎,支持 SQL/HiveQL 查询,DataFrame/Dataset 抽象,是目前最常用的组件。
- Spark Streaming:流计算引擎,处理实时数据流,支持 Kafka、Flume 等数据源。
- MLlib:分布式机器学习库,提供常用的分类、回归、聚类、推荐算法。
- GraphX:分布式图计算引擎,用于社交网络、路径计算等图场景。
1.4 官方学习资源指引
- 官方快速入门:新手入门第一站
- RDD 编程指南:核心编程规范
- Spark SQL 指南:结构化数据处理官方文档
- 官方示例仓库:官方提供的各类场景代码示例
二、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 任务执行完整流程
- 提交应用,Driver 启动,创建 SparkContext 上下文
- ApplicationMaster 向 YARN 申请 Executor 资源
- YARN 在 Core 节点上启动 Executor 进程
- Driver 将代码解析为 DAG,切分为 Stage 和 Task
- Task 分发到 Executor 上执行,结果返回 Driver
- 任务结束,释放资源
三、入门入口:Spark Shell 交互式开发
3.1 EMR 环境启动 Spark Shell
你的 EMR 集群已预装 Spark 3.5.3,无需任何配置,SSH 登录 Master 节点后直接执行启动命令:
# 启动交互式终端,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 上创建测试文件:
# 退出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
步骤2:Spark 实现词频统计
回到 spark-shell 中执行:
// 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)
步骤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 交互命令行
在终端中直接输入:
pyspark稍等几秒,出现 >>> 符号即表示进入了 Spark 环境。
3. 粘贴代码运行计算
将下面代码复制并粘贴到 >>> 后面,回车:
# 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 环境。

第三步:用 Spark 操作数据库(Spark SQL 实操)
Spark SQL 可以像写 SQL 语句一样操作 Hive 表。
1. 创建一个 Python 脚本文件
在终端输入:
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()
EOF2. 运行该脚本
python3 /root/demo_sql.py- 期望输出: 终端会显示一个格式化好的表格,包含
Alice和Bob的数据。
第四步:提交正式生产任务(spark-submit)
在实际工作中,我们不会一直在命令行里手动敲代码,而是将写好的代码用 spark-submit 提交给集群运行。
运行以下命令,将刚才的脚本提交给 YARN 资源管理器调度:
spark-submit --master yarn --deploy-mode client /root/demo_sql.py- 这个命令会把任务发送给您的从节点(
core-1-1和core-1-2)协同完成计算。

参数说明:
--master yarn:使用集群的 YARN 资源管理器,和交互式环境的运行架构一致--deploy-mode client:客户端模式,日志直接输出到当前终端,适合调试查看输出- 末尾填写脚本的绝对路径即可
总结:您只需要掌握这三个核心命令
pyspark:用来进入交互环境,边写代码边看结果(适合开发调试)。python3 xxx.py:用来在单机上快速运行写好的 Spark 脚本。spark-submit --master yarn xxx.py:用来将脚本提交到 Hadoop/YARN 集群真正进行分布式计算(适合生产环境)。
四、核心抽象:RDD 编程详解
4.1 RDD 核心概念
RDD(Resilient Distributed Dataset,弹性分布式数据集)是 Spark 最基础的抽象,代表一个不可变、可分区、里面的元素可并行计算的集合。
五大核心特性(官方定义):
- 分区(Partition):数据分片,是并行计算的基本单位
- 计算函数:每个分区都有一个计算函数
- 依赖关系:RDD 之间的血缘关系,记录转换过程
- 分区器:键值对 RDD 可指定分区策略
- 优先位置:计算任务调度到数据所在的节点,移动计算不移动数据
关键特性:
- 不可变性:RDD 一旦创建就不能修改,每次转换都会生成新的 RDD
- 弹性:数据丢失可以通过血缘关系重新计算,自动容错
- 分布式:数据分散存储在集群多个节点上
4.2 RDD 的三种创建方式
// 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,惰性执行,不会触发计算。
基础转换
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,是聚合计算的核心。
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)
行动算子触发真正的计算,返回结果或写入外部存储。
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,可以缓存到内存中,大幅提升性能。
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
// 从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 操作
// 查看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 中的表,无需任何额外配置。
// 查看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,常用语法和关系型数据库一致:
-- 建表
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 类型 | 地址 | 作用 | 生命周期 |
|---|---|---|---|
| 任务运行 UI | http://Master公网IP:4040 | 查看当前正在运行的任务详情 | 仅任务运行期间存在 |
| 历史服务 UI | http://Master公网IP:18080 | 查看所有历史任务记录 | 常驻服务,集群运行就存在 |
6.2 UI 核心功能
- Jobs 页面:查看所有作业,每个作业对应一个行动算子
- Stages 页面:查看阶段划分,DAG 可视化,了解任务执行流程
- Storage 页面:查看缓存的 RDD 存储情况、内存占用
- Executors 页面:查看所有 Executor 的资源使用、任务数、GC 情况
- Environment 页面:查看 Spark 配置参数、依赖 Jar 包
6.3 简单性能排查
- 任务慢:看 Stages 页面,找到耗时最长的 Stage,查看数据倾斜、GC 情况
- 内存不足:看 Executors 页面,查看堆内存使用,是否有 OOM 报错
- 任务失败:点击任务详情,查看错误日志,定位代码或数据问题
七、入门实战案例:访问日志分析
7.1 需求与数据准备
需求:模拟 Nginx 访问日志,统计访问量 Top3 的 IP、不同状态码的数量。
步骤1:创建测试数据
# 在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 实现
// 读取日志
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)
7.3 Spark SQL 实现
// 创建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()
八、学习路径与进阶方向
8.1 入门阶段(当前)
- 掌握 RDD 核心概念、常用算子
- 掌握 Spark SQL 基础查询、DataFrame 操作
- 理解 Spark 运行架构、惰性求值、DAG 执行原理
- 能独立完成简单的离线数据统计任务
8.2 进阶阶段
- 核心原理:DAG 划分、Stage 划分原理、宽窄依赖、共享变量
- 性能调优:资源配置、缓存策略、数据倾斜优化、Shuffle 调优
- 流计算:Spark Streaming / Structured Streaming 实时数据处理
- 生态整合:和 Kafka、HBase、Hive、S3 的整合使用
8.3 高级阶段
- Spark 内核原理、源码解析
- 企业级数仓建设、湖仓一体
- 机器学习 MLlib 实战
- 生产环境运维、监控、故障排查
九、使用注意事项(EMR 环境)
- 资源释放:学习结束后退出 spark-shell,释放集群资源;集群不用时及时释放,避免持续计费。
- 数据存储:重要数据存放在 HDFS 或 OSS,不要存在节点本地磁盘,释放集群会丢失。
- 提交方式:交互式开发用 spark-shell,正式任务用 spark-submit 提交脚本。
- 版本兼容:EMR 5.21.0 对应 Spark 3.5.3,代码和官方 Spark 3.x 完全兼容。