项目4:统计分析竞赛网站用户访问日志数据——Spark IDE编程
先修基础:项目1(Spark概述与集群搭建)、项目2(Scala基础)、项目3(Spark Shell编程)
目录
- 第一部分:项目背景与 Spark IDE 编程概述
- 第二部分:搭建 Spark 开发环境(IDEA + Scala 插件)
- 第三部分:AI 智能编程插件安装与使用
- 第四部分:编写第一个 Spark 程序
- 第五部分:在开发环境中运行 Spark 程序(本地模式)
- 第六部分:在集群环境中运行 Spark 程序(spark-submit)
- 第七部分:RDD 持久化(缓存)
- 第八部分:Spark 数据分区
- 第九部分:项目实战——竞赛网站用户访问日志分析
- 第十部分:常见问题与排错指南
- 第十一部分:实习 / 面试高频考点
- 附录:习题解析
第一部分:项目背景与 Spark IDE 编程概述
1.1 项目背景
为什么要学 Spark IDE 编程?
前面我们用 spark-shell 写代码,虽然方便,但有几个问题:
- 代码写一行执行一行,不能保存,关掉就没了
- 不能写复杂的多文件项目
- 不能调试
- 不能打包提交到集群运行
实际工作中,都是用 IDE(集成开发环境) 来写 Spark 程序的,最常用的就是 IntelliJ IDEA。
项目场景
某竞赛网站经过几年的运营,保存了众多用户对该网站的访问日志数据。为了研究用户的兴趣爱好并改善用户体验,需要对网站用户的访问日志数据进行分析。
数据文件:raceData.csv(2023年5月至2024年2月的用户访问日志)
本项目要做的事:
- 在 IDEA 中搭建 Spark 开发环境
- 编写 Spark 程序统计每月访问量
- 用自定义分区器按年份分区保存结果
- 打包提交到集群运行
1.2 spark-shell vs IDE 编程 对比
| 对比项 | spark-shell | IDEA 编程 |
|---|---|---|
| 使用方式 | 命令行交互式 | 编写完整程序 |
| 代码保存 | 不能保存 | 保存为文件,可复用 |
| SparkContext | 自动创建(sc变量) | 需要手动创建 |
| 适合场景 | 学习、调试、小数据测试 | 正式开发、生产环境 |
| 运行方式 | 直接在shell里跑 | 本地运行 / 打包提交集群 |
| 调试能力 | 弱 | 强(断点、单步等) |
| 项目管理 | 不支持 | 支持多文件、多模块 |
💡 理解:spark-shell 就像"草稿纸",IDEA 就像"正式作业本"。学习用 shell,开发用 IDEA。
第二部分:搭建 Spark 开发环境(IDEA + Scala 插件)
2.1 安装 IntelliJ IDEA
下载
- 官网下载 IntelliJ IDEA 安装包(社区版即可,免费)
- 版本:
ideaIC-2022.3.3.exe
安装步骤
- 双击安装包,进入安装向导
- 设置安装目录
- 完成一系列设置后,点击 Install
- 安装完成后重启电脑
首次启动
- 双击桌面图标启动
- 第一次启动问是否导入以前的设定 → 选"不导入"
- 欢迎界面可以切换主题(黑色/白色)
2.2 安装 Scala 插件
IDEA 默认不支持 Scala,需要安装插件。有两种安装方式:
方式一:在线安装(推荐,有网时用)
- 打开 IDEA,欢迎界面左侧选 Plugins
- 在 Marketplace 搜索框输入 Scala
- 点击 Scala 右侧的 Install 按钮
- 安装完成后点击 Restart IDE 重启
方式二:离线安装(没网时用)
- 提前下载好 Scala 插件包(如
scala-intellij-bin-2022.3.3.zip) - 在 Plugins 选项卡,点击齿轮图标 → 选 Install Plugin from Disk...
- 选择本地插件包路径
- 点击 OK 安装,然后重启 IDEA
⚠️ 注意:插件版本要和 IDEA 版本对应,否则可能装不上或用不了。
2.3 测试 Scala 插件
创建第一个 Scala 项目
- 欢迎界面选 New Project
- 填写项目信息:
- Name:
HelloWorld - Location:自定义存放路径
- Language:
Scala - Build System:
IntelliJ - JDK:
1.8.0_461 - Scala SDK:
2.13.16
- Name:
- 点击 Create 创建
项目结构
HelloWorld/
├── src/ # 源代码目录
├── .idea/ # IDEA配置文件(不用管)
└── HelloWorld.iml # 项目配置文件(不用管)编写 HelloWorld
- 右键 src → New → Package → 包名
com.tipdm.scalaDemo - 右键包名 → New → Scala Class → 类名
HelloWorld,类型选 Object - 编写代码:
package com.tipdm.scalaDemo
object HelloWorld {
def main(args: Array[String]): Unit = {
println("Hello World!")
}
}运行程序
右键代码空白处 → 选 Run 'HelloWorld'
控制台输出 Hello World! 就说明 Scala 环境配置成功了。
💡 说明:Scala 中
object是单例对象,main方法是程序入口,和 Java 的 main 方法类似。
第三部分:AI 智能编程插件安装与使用
3.1 常用 AI 编程插件
| 插件名 | 开发商 | 说明 |
|---|---|---|
| 通义灵码(Lingma) | 阿里云 | 国内常用,中文支持好 |
| TRAE AI | 字节跳动 | 原名 MarsCode,功能强大 |
💡 提示:一般装一个就够用了,也可以都装,根据需要切换。
3.2 安装通义灵码插件
- 打开 IDEA → Plugins → Marketplace
- 搜索 通义灵码(或英文 Lingma)
- 找到 "Lingma - Alibaba Cloud AI Coding Assistant"
- 点击 Install → 接受隐私协议 → 安装完成后重启 IDEA
TRAE AI 安装方式一样,搜索 TRAE 找到 "TRAE AI (formerly MarsCode): coding Assistant" 安装即可。
3.3 使用 AI 编程插件
第一步:登录
- 点击 IDEA 右下角的插件图标
- 按提示扫码或账号登录
- 登录成功后就可以用了
第二步:代码自动补全
- 写代码时会出现灰色的提示代码
- 按 Tab 键 确认使用提示
- 按 Esc 键 取消提示
第三步:对话式编程
- 右键代码空白处 → 选"通义灵码"或"Trae AI"
- 可以问问题、让AI写代码、解释代码、优化代码等
切换插件
- 如果装了多个 AI 插件,右侧侧边栏会有对应图标
- 点击不同图标切换使用
💡 给实习生的建议:
- AI 是辅助工具,不能完全依赖,要理解代码逻辑
- AI 生成的代码要自己检查,可能有 bug
- 用 AI 来提高效率,但核心知识还是要掌握
第四部分:编写第一个 Spark 程序
4.1 添加 Spark 开发依赖包
写 Spark 程序需要 Spark 的 jar 包,要先导入到项目中。
操作步骤
- 打开 Project Structure(快捷键:
Ctrl + Alt + Shift + S)- 或者菜单栏 File → Project Structure
- 选 Libraries 选项卡
- 点击 + 按钮 → 选 Java
- 找到本地 Spark 安装目录下的 jars 文件夹
- 点击 OK,Spark 依赖包就导入了
💡 说明:Spark 的 jars 文件夹里有所有需要的依赖包,一次性全部导入。
4.2 SparkContext 手动创建
在 spark-shell 中,sc 变量是自动创建好的。 但在 IDEA 中写程序,需要手动创建 SparkContext。
创建步骤
- 创建 SparkConf 对象:设置应用名称、运行模式等配置
- 创建 SparkContext 对象:传入 SparkConf 作为参数
- 用完后调用
sc.stop()关闭
代码模板
import org.apache.spark.{SparkConf, SparkContext}
object WordCount {
def main(args: Array[String]): Unit = {
// 1. 创建SparkConf(配置对象)
val conf = new SparkConf()
.setAppName("WordCount") // 设置应用名称
.setMaster("local") // 设置运行模式(本地模式)
// 2. 创建SparkContext(上下文对象,Spark程序入口)
val sc = new SparkContext(conf)
// 3. 写你的业务逻辑...
// ...
// 4. 关闭SparkContext
sc.stop()
}
}💡 理解:
- SparkConf 就像"配置清单",告诉 Spark 怎么运行
- SparkContext 就像"项目经理",是所有 Spark 操作的入口
- 程序结束要"下班",调用 stop() 关闭
4.3 第一个 Spark 程序:WordCount
package com.tipdm.scalaDemo
import org.apache.spark.{SparkConf, SparkContext}
object WordCount {
def main(args: Array[String]): Unit = {
// 创建SparkConf和SparkContext
val conf = new SparkConf().setAppName("WordCount").setMaster("local")
val sc = new SparkContext(conf)
// 设置Hadoop安装路径(Windows本地运行需要)
System.setProperty("hadoop.home.dir", "E:\\hadoop-3.3.6")
// 单词计数:读取文件 → 切分单词 → 映射为(单词,1) → 按Key求和
val count = sc.textFile("E:\\data\\words.txt")
.flatMap(x => x.split(" "))
.map(x => (x, 1))
.reduceByKey((x, y) => x + y)
// 打印结果
count.foreach(println)
// 关闭
sc.stop()
}
}代码说明
| 代码 | 说明 |
|---|---|
new SparkConf() | 创建配置对象 |
.setAppName("WordCount") | 设置应用名(在Spark UI上显示) |
.setMaster("local") | 设置本地模式运行 |
new SparkContext(conf) | 创建Spark上下文 |
System.setProperty(...) | 设置Hadoop路径(Windows需要) |
sc.textFile(...) | 读取文件创建RDD |
.flatMap(...).map(...).reduceByKey(...) | 单词计数逻辑 |
sc.stop() | 关闭Spark上下文 |
第五部分:在开发环境中运行 Spark 程序(本地模式)
在 IDEA 中直接运行 Spark 程序,需要设置三个东西:
- 运行模式(必须设置,否则找不到 master)
- Hadoop bin 路径(Windows 需要)
- 自定义输入参数(如果程序需要参数)
5.1 设置运行模式
方式一:在代码中设置(推荐,学习时用)
val conf = new SparkConf().setAppName("WordCount").setMaster("local")方式二:在 IDEA 运行配置中设置
- 菜单栏 Run → Edit Configurations...
- 找到 VM options 输入框
- 输入:
-Dspark.master=local - 点击 OK
💡 两种方式的区别:
- 代码中设置:写死了,打包到集群运行时要改代码
- VM options 设置:不写死在代码里,更灵活
- 实际开发:代码里不写 setMaster,打包后通过 spark-submit 指定
5.2 指定 Hadoop bin 路径(Windows 特有)
在 Windows 本地运行 Spark 程序,需要指定 Hadoop 的 bin 文件夹路径。
方式一:在代码中设置
System.setProperty("hadoop.home.dir", "E:\\hadoop-3.3.6")方式二:设置 Windows 环境变量
- 在 Path 环境变量中添加 Hadoop 的 bin 文件夹路径
- 这样程序运行时就不用每次指定了
还需要两个文件
Hadoop 的 bin 文件夹里需要有:
winutils.exehadoop.dll
这两个文件可以从 GitHub 上下载对应版本的。
💡 说明:Linux 环境下不需要设置这个,只有 Windows 本地开发时才需要。
5.3 设置自定义输入参数
如果程序需要从外部传入参数(比如输入路径、输出路径),就要设置 Program arguments。
代码中接收参数
val input = args(0) // 第一个参数:输入路径
val output = args(1) // 第二个参数:输出路径IDEA 中设置参数
- Run → Edit Configurations...
- 找到 Program arguments 输入框
- 输入参数(空格分隔),比如:
E:\data\words.txt E:\data\output - 点击 OK
5.4 本地运行完整示例
代码(带参数版本)
package com.tipdm.scalaDemo
import org.apache.spark.{SparkConf, SparkContext}
object WordCount {
def main(args: Array[String]): Unit = {
val conf = new SparkConf().setAppName("WordCount").setMaster("local")
val sc = new SparkContext(conf)
System.setProperty("hadoop.home.dir", "E:\\hadoop-3.3.6")
// 从参数获取输入输出路径
val input = args(0)
val output = args(1)
// 单词计数
val count = sc.textFile(input)
.flatMap(x => x.split(" "))
.map(x => (x, 1))
.reduceByKey((x, y) => x + y)
// 保存结果
count.repartition(1).saveAsTextFile(output)
sc.stop()
}
}设置参数
Program arguments 填:E:\data\words.txt E:\data\wordcount_out
运行
右键 → Run,等程序跑完,去输出路径看结果。
第六部分:在集群环境中运行 Spark 程序(spark-submit)
本地模式只适合小数据测试,生产环境都是打包成 JAR 包,提交到集群运行。
6.1 整体流程
IDEA写代码 → 打包成JAR → 上传到Linux → spark-submit提交到集群 → 查看结果6.2 修改代码(去掉本地模式设置)
打包到集群运行的代码,不要在代码里写死 setMaster("local"),让 spark-submit 来指定。
package com.tipdm.scalaDemo
import org.apache.spark.{SparkConf, SparkContext}
object WordCount {
def main(args: Array[String]): Unit = {
// 注意:不要写 setMaster("local")!由spark-submit指定
val conf = new SparkConf().setAppName("WordCount")
val sc = new SparkContext(conf)
val input = args(0)
val output = args(1)
val count = sc.textFile(input)
.flatMap(x => x.split(" "))
.map(x => (x, 1))
.reduceByKey((x, y) => x + y)
count.repartition(1).saveAsTextFile(output)
sc.stop()
}
}⚠️ 重要:集群运行的代码里不要 setMaster("local"),也不要 System.setProperty("hadoop.home.dir"),这些都是 Windows 本地开发才需要的。
6.3 在 IDEA 中打包工程(生成 JAR 包)
步骤一:配置 Artifact
- File → Project Structure(Ctrl+Alt+Shift+S)
- 选 Artifacts 选项卡
- 点击 + → JAR → Empty
- Name 填 JAR 包名称(如
word) - 右侧找到你的项目 → 双击
'项目名' compile output- 它会移到左侧,表示已添加到 JAR 包中
- 点击 OK
步骤二:编译生成 JAR
- 菜单栏 Build → Build Artifacts...
- 选择你的 JAR 包名 → Build
- 等待编译完成
步骤三:找到 JAR 包
- 生成的 JAR 包在工程目录的
out/artifacts/目录下 - 右键 JAR 包 → Open In → Explorer 可以打开文件位置
6.4 spark-submit 提交命令
基本格式
spark-submit \
--master <master-url> \
--deploy-mode <deploy-mode> \
--class <main-class> \
<application-jar> \
[application-arguments]常用参数说明
| 参数 | 说明 | 示例 |
|---|---|---|
--master | 集群地址(运行模式) | spark://master:7077 / yarn |
--deploy-mode | 部署模式:client / cluster | cluster |
--class | 主类入口(包名+类名) | com.tipdm.scalaDemo.WordCount |
--conf | 任意 Spark 配置属性 | --conf spark.executor.memory=1g |
--executor-memory | 每个 Executor 的内存 | 1G |
--executor-cores | 每个 Executor 的 CPU 核数 | 1 |
--num-executors | Executor 数量 | 2 |
application-jar | JAR 包路径 | /opt/data/word.jar |
application-arguments | 传给 main 方法的参数 | 输入路径 输出路径 |
示例1:Standalone 模式提交
cd $SPARK_HOME/bin
spark-submit \
--master spark://master:7077 \
--class com.tipdm.scalaDemo.WordCount \
/opt/data/word.jar \
/tipdm/data/words.txt \
/tipdm/data/wordcount示例2:YARN cluster 模式提交
cd $SPARK_HOME/bin
spark-submit \
--master yarn \
--deploy-mode cluster \
--class com.tipdm.scalaDemo.WordCount \
/opt/data/word.jar \
/tipdm/data/words.txt \
/tipdm/data/wordcount示例3:带资源配置的提交
spark-submit \
--master yarn \
--deploy-mode cluster \
--executor-memory 1G \
--executor-cores 1 \
--class com.tipdm.scalaDemo.WordCount \
/opt/data/word.jar \
/tipdm/data/words.txt \
/tipdm/data/wordcount26.5 配置优先级(面试常问)
Spark 配置可以在多个地方设置,优先级从高到低:
| 优先级 | 配置位置 | 说明 |
|---|---|---|
| ⭐⭐⭐⭐⭐ 最高 | 代码中 conf.set() | 写死在代码里的 |
| ⭐⭐⭐⭐ | spark-submit 命令行参数 | --conf、--executor-memory 等 |
| ⭐⭐⭐ | 配置文件 spark-defaults.conf | Spark 安装目录 conf 下 |
| ⭐⭐ 最低 | 系统默认值 | Spark 内置默认值 |
💡 记忆口诀:代码 > 命令行 > 配置文件 > 默认值
越"近"的优先级越高,代码里写死的最高,命令行其次,配置文件再次,默认值最低。
6.6 提交到集群运行完整流程
步骤
- 启动集群:启动 Hadoop 和 Spark
- 上传 JAR 包:把 Windows 的 JAR 包传到 Linux 的
/opt/data/目录 - 上传数据:把数据文件上传到 HDFS
- 提交运行:用 spark-submit 提交
- 查看结果:在 HDFS 上查看输出
完整命令示例
# 1. 启动Hadoop和Spark(已启动就跳过)
start-dfs.sh
start-yarn.sh
$SPARK_HOME/sbin/start-all.sh
# 2. 上传数据到HDFS
hdfs dfs -put /opt/data/words.txt /tipdm/data/
# 3. 提交Spark程序
cd $SPARK_HOME/bin
spark-submit \
--master yarn --deploy-mode cluster \
--executor-memory 1G --executor-cores 1 \
--class com.tipdm.scalaDemo.WordCount \
/opt/data/word.jar \
/tipdm/data/words.txt \
/tipdm/data/wordcount
# 4. 查看结果
hdfs dfs -cat /tipdm/data/wordcount/part-00000第七部分:RDD 持久化(缓存)
7.1 为什么需要持久化?
因为 Spark RDD 是惰性求值的,每次调用行动算子都会从头重新计算。
如果一个 RDD 要被多次使用,每次都重新算就太浪费了。 持久化(缓存)就是把 RDD 的计算结果保存到内存/磁盘,下次用直接取,不用重算。
举个例子
val data = sc.parallelize(List(1,2,3,4,5,6)).map(x => x*x)
// 第一次用:从头计算(map要执行)
println(data.sum())
// 第二次用:又从头计算(map又要执行一遍!)
println(data.mean())如果 data 持久化了,第二次就直接从内存取,不用再算 map 了。
💡 形象比喻:
- 不持久化:每次做题都从翻书找公式开始
- 持久化:把公式抄在草稿纸上,下次直接看草稿纸
7.2 持久化方法
两个方法
| 方法 | 说明 |
|---|---|
cache() | 默认存储级别(MEMORY_ONLY),快捷方法 |
persist(存储级别) | 可指定存储级别,更灵活 |
代码示例
import org.apache.spark.storage.StorageLevel
val data = sc.parallelize(List(1, 2, 3, 4, 5, 6)).map(x => x*x)
// 持久化:存储到内存
data.persist(StorageLevel.MEMORY_ONLY)
// 第一次使用:触发计算,并把结果存到内存
println(data.sum())
// 第二次使用:直接从内存取,不用重算
println(data.mean())💡 注意:persist() 是转换算子,不会立即触发持久化,要等第一次行动算子执行时才会真正持久化。
7.3 存储级别
org.apache.spark.storage.StorageLevel 中的常用存储级别:
| 存储级别 | 说明 | 内存 | 磁盘 | 序列化 | 副本数 |
|---|---|---|---|---|---|
| MEMORY_ONLY | 只存内存(反序列化) | ✅ | ❌ | ❌ | 1 |
| MEMORY_AND_DISK | 内存不够放磁盘 | ✅ | ✅ | ❌ | 1 |
| MEMORY_ONLY_SER | 只存内存(序列化) | ✅ | ❌ | ✅ | 1 |
| MEMORY_AND_DISK_SER | 内存+磁盘(序列化) | ✅ | ✅ | ✅ | 1 |
| DISK_ONLY | 只存磁盘 | ❌ | ✅ | - | 1 |
| MEMORY_ONLY_2 | 只存内存,2副本 | ✅ | ❌ | ❌ | 2 |
💡 说明:
- 序列化:把对象转成字节数组,更省空间,但读取时要反序列化,费点 CPU
- 副本:存多份,防止一个节点挂了数据丢失,但更占空间
- 末尾加
_2表示 2 副本,如MEMORY_ONLY_2
7.4 持久化注意事项
① 存储级别只能设置一次
- 只有未曾设置存储级别的 RDD 才能设置
- 已经设置了的不能修改
- 如果要改,得先
unpersist()再重新设置
② 内存不足时的策略
- 内存不够时,用 LRU(最近最少使用) 策略清除最老的分区
- 内存中的缓存不是永久的,可能被清掉
- 磁盘上的缓存不会被清
③ 手动清除缓存
// 手动清除不需要的缓存数据
data.unpersist()💡 什么时候用持久化?
- 一个 RDD 被多次使用时(比如迭代计算、机器学习)
- 计算代价很高的 RDD(计算一次很费时间)
- 数据量不大,能放进内存的
第八部分:Spark 数据分区
8.1 为什么要控制数据分区?
Spark RDD 是由多个分区组成的,数据分布在多个节点上。 分布式计算中,网络通信的代价很大。 控制数据分区、减少网络传输(Shuffle),是提升性能的重要手段。
分区的两个维度
- 分区个数:有多少个分区
- 分区方式:数据怎么分配到各个分区(分区器)
💡 理解:分区就像"分箱子装东西",分区个数是"有几个箱子",分区方式是"按什么规则往箱子里放"。
8.2 分区器(Partitioner)
重要前提
只有键值对 RDD 才能设置分区方式(分区器) 非键值对 RDD 的分区方式是 None(只能设置分区个数)。
Spark 内置的两种分区器
| 分区器 | 说明 | 特点 |
|---|---|---|
| HashPartitioner | 哈希分区 | 根据 Key 的哈希值取模分区,简单常用 |
| RangePartitioner | 范围分区 | 把一定范围的数据放到一个分区,排序后均匀分布 |
查看分区方式
// 查看RDD的分区器
rdd.partitioner // 返回Option[Partitioner]
// 判断是否有分区器
rdd.partitioner.isDefined
// 获取具体的分区器
rdd.partitioner.get8.3 partitionBy() 方法
设置分区方式用 partitionBy() 方法,传入一个分区器。
// 对键值对RDD设置哈希分区,3个分区
val partitioned = kv_rdd.partitionBy(new HashPartitioner(3))
// 查看分区数
partitioned.partitions.size8.4 自定义分区器(重点+难点)
内置分区器满足不了需求时,可以自定义分区器。
实现步骤
继承 org.apache.spark.Partitioner 类,实现 3 个方法:
| 方法 | 签名 | 作用 |
|---|---|---|
numPartitions | Int | 返回分区个数 |
getPartition(key: Any) | Int | 根据 Key 返回分区 ID(0 到 numPartitions-1) |
equals(other: Any) | Boolean | 判断两个分区器是否相等 |
示例:按奇偶性分区
自定义一个分区器,把偶数 Key 放分区 0,奇数 Key 放分区 1。
import org.apache.spark.Partitioner
class MyPartition(numParts: Int) extends Partitioner {
// 1. 返回分区个数
override def numPartitions: Int = numParts
// 2. 根据Key返回分区ID
override def getPartition(key: Any): Int = {
if (key.toString.toInt % 2 == 0) {
return 0 // 偶数 → 分区0
} else {
return 1 // 奇数 → 分区1
}
}
// 3. equals方法(判断两个分区器是否一样)
override def equals(other: Any): Boolean = other match {
case mypartition: MyPartition =>
mypartition.numPartitions == numPartitions
case _ => false
}
}使用自定义分区器
import org.apache.spark.{SparkConf, SparkContext}
object ToDistribute {
def main(args: Array[String]): Unit = {
val conf = new SparkConf().setAppName("PartitionTest")
val sc = new SparkContext(conf)
val input = args(0)
val output = args(1)
// 读取数据,创建键值对RDD(第1列是Key,第2列是Value)
val data = sc.textFile(input).map { x =>
val y = x.split(",")
(y(0), y(1))
}
// 使用自定义分区器
val data2 = data.partitionBy(new MyPartition(2))
// 保存结果(会生成2个文件,part-00000和part-00001)
data2.saveAsTextFile(output)
sc.stop()
}
}提交运行
spark-submit \
--master spark://master:7077 \
--executor-memory 512M --executor-cores 1 \
--class com.tipdm.scalaDemo.ToDistribute \
/opt/data/word.jar \
/tipdm/data/user.txt \
/tipdm/data/part1💡 结果说明:
- part-00000:偶数 Key 的数据
- part-00001:奇数 Key 的数据
8.5 coalesce() 和 repartition() 方法
除了自定义分区器,还有两个简单的重分区方法,任何类型的 RDD 都能用。
① coalesce() 方法
语法:
coalesce(numPartitions: Int, shuffle: Boolean = false)参数说明:
numPartitions:重分区后的分区数shuffle:是否进行 Shuffle,默认 false
规则:
- shuffle = false(默认):只能减少分区数,不能增加(增加的话不生效)
- shuffle = true:可以增加也可以减少分区数(会触发 Shuffle)
② repartition() 方法
语法:
repartition(numPartitions: Int)本质:就是 coalesce(numPartitions, shuffle = true) 的简写。 特点:一定会触发 Shuffle,可以增加也可以减少分区数。
代码示例
val rdd = sc.parallelize(List(1, 2, 3, 4, 5, 6, 7))
// 查看原始分区数
rdd.partitions.size
// coalesce减少分区(不shuffle)
rdd.coalesce(5).partitions.size // 变成5个
// coalesce增加分区(不shuffle → 不生效,分区数不变)
rdd.coalesce(10).partitions.size // 还是原来的数量
// coalesce增加分区(shuffle=true → 生效)
rdd.coalesce(10, true).partitions.size // 变成10个
// repartition(等价于coalesce(10, true))
rdd.repartition(10).partitions.size // 变成10个coalesce vs repartition 对比
| 对比项 | coalesce(shuffle=false) | coalesce(shuffle=true) | repartition |
|---|---|---|---|
| 能否增加分区 | ❌ 不能 | ✅ 能 | ✅ 能 |
| 能否减少分区 | ✅ 能 | ✅ 能 | ✅ 能 |
| 是否触发Shuffle | ❌ 不触发 | ✅ 触发 | ✅ 触发 |
| 性能 | 好(无Shuffle) | 差(有Shuffle) | 差(有Shuffle) |
| 适用场景 | 合并小分区、减少分区数 | 增加分区数、均匀分布 | 简单重分区 |
💡 使用建议:
- 只是想减少分区数 → 用 coalesce(不触发 Shuffle,性能好)
- 想增加分区数,或者让数据更均匀 → 用 repartition
第九部分:项目实战——竞赛网站用户访问日志分析
9.0 数据说明
数据文件
raceData.csv:竞赛网站用户访问日志数据(2023年5月至2024年2月)
时间字段格式
- 2023年数据:
yyyy/MM/dd HH:mm:ss(斜杠分隔) - 2024年数据:
yyyy-MM-dd HH:mm:ss(横杠分隔)
⚠️ 注意:两年的日期格式不一样,需要统一处理。
字段说明
数据有多个字段,第 5 列(索引 4)是访问时间。
任务4.1:计算竞赛网站每月的访问量
需求
统计每个月的网站访问量。
实现思路
- 读取原始数据
- 取出时间字段(第5列,索引4)
- 统一日期格式,提取"年-月"
- 映射为 (年-月, 1) 的键值对
- 按 Key 求和 → 得到每月访问量
实现代码
import org.apache.spark.{SparkConf, SparkContext}
object race {
def main(args: Array[String]): Unit = {
// 1. 创建SparkContext
val conf = new SparkConf().setAppName("race")
val sc = new SparkContext(conf)
// 2. 获取输入参数
val input = args(0)
// 3. 读取数据创建RDD
val rdd = sc.textFile(input)
// 4. 处理数据,统计每月访问量
val result = rdd
.map(x => x.split(",")(4)) // 取第5列(时间)
.map(y => y.split(" ")(0).replace("/", "-")) // 取日期部分,把/换成-
.map(z => (z.split("-")(0) + "-" + z.split("-")(1), 1)) // 提取"年-月",值为1
.reduceByKey((x, y) => x + y) // 按月求和
// 5. 打印结果
result.collect().foreach(println)
// 6. 关闭
sc.stop()
}
}代码逐行解析
| 代码 | 说明 |
|---|---|
x.split(",")(4) | 按逗号分割,取第5个元素(索引4)= 时间字段 |
y.split(" ")(0) | 按空格分割,取第一部分 = 日期部分 |
.replace("/", "-") | 把斜杠换成横杠,统一格式 |
z.split("-")(0) + "-" + z.split("-")(1) | 取年和月,拼成"年-月" |
.reduceByKey(_ + _) | 按月分组求和 |
提交运行
cd $SPARK_HOME/bin
spark-submit \
--master spark://master:7077 \
--class race \
/opt/data/race.jar \
/tipdm/data/raceData.csv任务4.2:自定义分区保存结果
需求
把每月访问量的结果按年份分区保存:
- 2023年的数据 → 分区 0(part-00000)
- 2024年的数据 → 分区 1(part-00001)
实现思路
- 先统计每月访问量(和任务4.1一样)
- 自定义分区器,按年份分区
- 用 partitionBy 应用分区器
- 保存结果到 HDFS
实现代码
import org.apache.spark.{Partitioner, SparkConf, SparkContext}
object race {
def main(args: Array[String]): Unit = {
// 1. 创建SparkContext
val conf = new SparkConf().setAppName("race")
val sc = new SparkContext(conf)
// 2. 获取参数
val input = args(0)
val output = args(1)
// 3. 读取数据
val rdd = sc.textFile(input)
// 4. 统计每月访问量
val result = rdd
.map(x => x.split(",")(4))
.map(y => y.split(" ")(0).replace("/", "-"))
.map(z => (z.split("-")(0) + "-" + z.split("-")(1), 1))
.reduceByKey((x, y) => x + y)
// 5. 使用自定义分区器分区,并保存
result.partitionBy(new MyPartitioner).saveAsTextFile(output)
// 6. 关闭
sc.stop()
}
// 自定义分区器:按年份分区
class MyPartitioner extends Partitioner {
// 分区数:2个(2023年一个,2024年一个)
override def numPartitions: Int = 2
// 根据Key(年-月)决定分到哪个分区
override def getPartition(key: Any): Int = {
key match {
case "2024-1" => 1 // 2024年1月 → 分区1
case "2024-2" => 1 // 2024年2月 → 分区1
case _ => 0 // 其他(都是2023年的)→ 分区0
}
}
}
}代码说明
自定义分区器 MyPartitioner
numPartitions = 2:两个分区getPartition:根据 Key("年-月"格式)判断- 2024 年的 → 分区 1
- 其他(2023 年的)→ 分区 0
💡 优化思路:这里用 match 写死了,如果年份更多,可以提取年份后判断:
scalaoverride def getPartition(key: Any): Int = { val year = key.toString.split("-")(0).toInt if (year == 2024) 1 else 0 }
提交运行
spark-submit \
--master spark://master:7077 \
--class race \
/opt/data/race.jar \
/tipdm/data/raceData.csv \
/tipdm/data/race查看结果
# 查看输出目录
hdfs dfs -ls /tipdm/data/race
# 查看分区0(2023年的数据)
hdfs dfs -cat /tipdm/data/race/part-00000
# 查看分区1(2024年的数据)
hdfs dfs -cat /tipdm/data/race/part-00001第十部分:常见问题与排错指南
10.1 环境配置类问题
问题1:IDEA 中找不到 Scala SDK
现象:创建 Scala 项目时没有 Scala SDK 选项 原因:Scala 插件没装好,或者 SDK 没配置 解决:
- 确认 Scala 插件已安装并启用
- 在 Project Structure → Global Libraries 中添加 Scala SDK
问题2:运行时报错找不到 hadoop.home.dir
现象:Windows 下运行报 java.io.IOException: Could not locate executable null\bin\winutils.exe原因:没有设置 Hadoop 路径 解决:
- 代码中加:
System.setProperty("hadoop.home.dir", "你的hadoop路径") - 或者设置 Windows 环境变量
- 确保 hadoop 的 bin 目录下有 winutils.exe 和 hadoop.dll
问题3:运行时报错 "A master URL must be set"
现象:报错 org.apache.spark.SparkException: A master URL must be set in your configuration原因:没有设置运行模式(master URL) 解决:
- 代码中加
.setMaster("local") - 或者在 VM options 中加
-Dspark.master=local
10.2 打包提交类问题
问题4:打包后运行报 ClassNotFoundException
现象:spark-submit 提交后报找不到主类 原因:
- --class 后面的类名写错了(包名+类名要写全)
- JAR 包里没有包含编译好的 class 文件 解决:
- 确认类名写对了(要写全包名,如
com.tipdm.scalaDemo.WordCount) - 检查 Artifact 配置中是否添加了 compile output
问题5:提交到 YARN 集群报错
现象:报各种 YARN 相关错误 原因:
- Hadoop 没启动
- 环境变量没配置好(HADOOP_CONF_DIR、YARN_CONF_DIR)
- Spark 配置文件没改对 解决:
- 确认 Hadoop 正常启动
- 检查 spark-env.sh 中的 HADOOP_CONF_DIR 配置
- 检查 spark-defaults.conf 配置
问题6:Executor 内存不足
现象:Executor 挂掉,报 OOM 或 GC overhead 原因:分配的 Executor 内存不够 解决:
- 调大
--executor-memory - 增加分区数,减少每个分区的数据量
- 优化代码,避免数据倾斜
10.3 持久化和分区类问题
问题7:persist() 后还是重新计算了
原因:
- persist 是转换算子,第一次行动算子才会真正持久化
- 内存不够,部分分区被淘汰了
- 用了 MEMORY_ONLY,内存不够时直接丢了 解决:
- 确认在第一次行动算子之后才会有缓存
- 内存不够就用 MEMORY_AND_DISK
- 数据太大就不要全缓存了
问题8:自定义分区器不生效
现象:分区结果和预期不一样 原因:
- getPartition 方法逻辑写错了
- 返回的分区 ID 超出范围(必须是 0 到 numPartitions-1)
- RDD 不是键值对类型 解决:
- 检查 getPartition 的逻辑
- 确保分区 ID 在合法范围内
- 确认是键值对 RDD 才能用 partitionBy
问题9:coalesce 增加分区数不生效
原因:coalesce 默认 shuffle=false,只能减少分区数,不能增加 解决:
- 要增加分区数,用
coalesce(n, shuffle = true) - 或者直接用
repartition(n)
10.4 排错通用思路
- 看报错信息:先看异常类型和错误消息
- 本地先测通:先在本地模式(local)跑通,再打包提交集群
- 看 Spark UI:4040 端口看 Job、Stage、Executor 情况
- 看日志:YARN 日志或 Worker 日志找具体错误
- 小数据测试:先用小数据验证逻辑
第十一部分:实习 / 面试高频考点
11.1 概念类(高频)
Q1:spark-shell 和 IDEA 编程的区别?
spark-shell 是交互式命令行,自动创建 SparkContext,适合学习调试,代码不能保存 IDEA 编程是写完整程序,手动创建 SparkContext,适合正式开发,代码可保存、可打包、可调试 实际工作中都用 IDEA 开发,打包提交集群运行
Q2:Spark 程序的入口是什么?怎么创建?
SparkContext 是 Spark 程序的入口 创建步骤:先创建 SparkConf(设置应用名、运行模式等),再 new SparkContext(conf) spark-shell 中自动创建了 sc 变量,IDEA 中要手动创建
Q3:spark-submit 常用参数有哪些?
--master:指定集群地址(运行模式) --deploy-mode:部署模式(client/cluster) --class:主类(包名+类名) --executor-memory:每个 Executor 内存 --executor-cores:每个 Executor CPU 核数 --num-executors:Executor 数量 --conf:任意 Spark 配置
Q4:Spark 配置的优先级?
从高到低:代码中 set() > spark-submit 命令行 > 配置文件 spark-defaults.conf > 系统默认值 越"近"的优先级越高
Q5:为什么要持久化 RDD?什么时候用?
因为 RDD 是惰性求值的,每次行动算子都从头计算 如果一个 RDD 被多次使用,每次都重算很浪费 持久化把结果存到内存/磁盘,下次直接用,不用重算 适用场景:迭代计算、机器学习、一个 RDD 多次使用的场景
Q6:cache() 和 persist() 的区别?
cache() 是 persist() 的快捷方式,默认用 MEMORY_ONLY 存储级别 persist() 可以指定各种存储级别,更灵活 cache() = persist(StorageLevel.MEMORY_ONLY)
Q7:常用的存储级别有哪些?
MEMORY_ONLY:只存内存,反序列化,默认 MEMORY_AND_DISK:内存不够放磁盘 MEMORY_ONLY_SER:只存内存,序列化(更省空间) DISK_ONLY:只存磁盘 末尾加 _2 表示 2 副本
Q8:什么是分区器?有哪些内置分区器?
分区器决定数据怎么分配到各个分区 只有键值对 RDD 才能设置分区器 内置两种:
- HashPartitioner:按 Key 哈希值取模分区
- RangePartitioner:按范围分区,排序后均匀分布
Q9:coalesce 和 repartition 的区别?
coalesce:默认不触发 Shuffle,只能减少分区数,性能好 repartition:一定会触发 Shuffle,可以增也可以减,数据更均匀 本质:repartition = coalesce(shuffle=true) 使用建议:减少分区用 coalesce,增加分区用 repartition
Q10:怎么自定义分区器?
继承 org.apache.spark.Partitioner 类 实现三个方法:
- numPartitions:返回分区个数
- getPartition(key):根据Key返回分区ID(0到numPartitions-1)
- equals(other):判断两个分区器是否相等
11.2 原理类(中频)
Q11:持久化的数据丢失了怎么办?
如果用的是内存存储,内存不够时会用 LRU 策略淘汰老的分区 丢失的分区下次使用时会自动重新计算(RDD 的容错机制) 如果不想丢失,可以用磁盘存储,或者用副本(_2后缀)
Q12:为什么 HashPartitioner 可能导致数据倾斜?
哈希分区是按 Key 的哈希值取模,如果某些 Key 特别多(热点 Key),就会导致某个分区数据特别多 数据倾斜会导致某些 Task 特别慢,整体性能差 解决方法:自定义分区器、加盐、增加分区数等
Q13:client 模式和 cluster 模式的区别?
client 模式:Driver 运行在提交客户端,能直接看到日志,适合调试 cluster 模式:Driver 运行在集群中(AM 里),适合生产环境 YARN 模式下常用 cluster 模式
11.3 实操类(高频)
Q14:怎么把 IDEA 项目打包成 JAR?
- File → Project Structure → Artifacts → + → JAR → Empty
- 设置名称,添加 compile output
- Build → Build Artifacts → Build
- 去 out/artifacts 目录下找 JAR 包
Q15:Windows 本地运行 Spark 程序需要注意什么?
- 要设置 setMaster("local")
- 要设置 hadoop.home.dir
- Hadoop 的 bin 目录下要有 winutils.exe 和 hadoop.dll
- 只是本地测试用,生产环境代码里不要写死这些
Q16:怎么查看 Spark 程序运行情况?
- Spark Web UI:4040 端口(本地运行时)
- Standalone 集群:8080 端口看 Master UI
- YARN 集群:8088 端口看 ResourceManager UI
- 看日志:YARN 日志或 Worker 日志
附录:习题解析
选择题解析
1、答案:A 解析:使用 spark-submit 命令提交运行程序时,--master 用于指定运行模式,--class 用于指定主类,--name 用于指定应用名。
2、答案:B 解析:使用 spark-submit 命令提交运行程序时,--jars 用于为应用程序添加运行时所需的依赖 JAR 包,--deploy-mode 用于指定部署模式(如 YARN 集群模式和 YARN 客户端模式),--executor-memory 用于指定每个 Executor 可使用的内存大小。
3、答案:A 解析:使用 spark-submit 命令提交运行程序时,--name 用于指定应用名。
4、答案:A 解析:只有键值对 RDD 才能设置分区方式(分区器),非键值对 RDD 只能设置分区个数。
5、答案:B 解析:contains() 方法需接收字符串或字符参数,若未加引号会被视为变量;Scala 编程中定义列表时应使用圆括号 (),而非花括号 {}。
6、答案:D 解析:在 RDD 的持久化操作中,cache() 方法是使用默认存储级别的快捷方法,即 StorageLevel.MEMORY_ONLY(将反序列化的对象存入内存)。
7、答案:D 解析:SparkContext 是 Spark 应用程序的上下文和入口,任何 Spark 程序都从 SparkContext 对象开始。
8、答案:D 解析:Spark 官方并未提供 MyPartitioner 类,实现自定义分区器需要继承 org.apache.spark.Partitioner 类并自己实现。
9、答案:B 解析:使用 spark-submit 命令提交运行程序时,--master 用于指定运行模式,yarn-cluster 不是有效的运行模式,正确写法应为 --master yarn,并通过 --deploy-mode cluster 指定部署模式。
10、答案:C 解析:所有选项中均是先将每个字符串按空格分割为单词并转换为 (单词, 1) 的键值对,A、B 选项中后续操作为对相同键的值求和,统计单词数量;D 选项后续操作为将相同键的值分组为 Iterable 集合,并计算每个 Iterable 集合的元素数量,即单词出现次数。A、B、D 均可实现单词计数功能。
操作题解析
题目
读取 house.txt 文件,统计各城区的二手房套数,并将结果保存到 HDFS。
数据说明
- house.txt:二手房数据,制表符(\t)分隔
- 第 4 列(索引 3):城区
- 第 7 列(索引 6):套数
- 第 6 列(索引 5):小区名称(用于过滤空值)
实现代码
import org.apache.spark.{SparkConf, SparkContext}
object houseCount {
def main(args: Array[String]): Unit = {
// 创建SparkContext
val conf = new SparkConf().setAppName("houseCount")
val sc = new SparkContext(conf)
// 设置程序参数
val input = args(0)
val output = args(1)
// 读取数据
val data = sc.textFile(input)
// 数据清洗:过滤掉第6个字段为空的行
val dataClean = data.filter(_.split("\\t")(5) != "")
.map(line => {
val data = line.split("\\t")
(data)
})
// 缓存到内存
dataClean.cache()
// 统计各城区的二手房套数
val peopleCount = dataClean
.map(data => (data(3), data(6).toInt))
.reduceByKey(_ + _)
// 保存结果到HDFS
peopleCount.repartition(1).saveAsTextFile(output)
}
}提交运行
# 上传数据到HDFS
hdfs dfs -put /opt/data/house.txt /tipdm/data
# 打包后提交到YARN集群运行
cd $SPARK_HOME/bin
spark-submit --master yarn --deploy-mode cluster \
--class houseCount /opt/data/house.jar \
/tipdm/data/house.txt \
/tipdm/data/houseResult
# 查看HDFS结果
hdfs dfs -cat /tipdm/data/houseResult/part-00000代码说明
| 代码 | 说明 |
|---|---|
_.split("\\t") | 按制表符分割(注意转义) |
filter(_.split("\\t")(5) != "") | 过滤第6列为空的行(数据清洗) |
dataClean.cache() | 缓存清洗后的数据,多次使用时不用重算 |
(data(3), data(6).toInt) | 取城区(第4列)和套数(第7列),套数转Int |
reduceByKey(_ + _) | 按城区求和 |
.repartition(1) | 合并成一个分区,输出一个文件 |
笔记版本:V1.0 对应教材:《Spark大数据技术与应用(第3版)》人民邮电出版社 对应项目:项目4 统计分析竞赛网站用户访问日志数据——Spark IDE编程 最后更新:2026年8月