Skip to content

项目4:统计分析竞赛网站用户访问日志数据——Spark IDE编程 ​

先修基础:项目1(Spark概述与集群搭建)、项目2(Scala基础)、项目3(Spark Shell编程)


目录 ​


第一部分:项目背景与 Spark IDE 编程概述 ​

1.1 项目背景 ​

为什么要学 Spark IDE 编程? ​

前面我们用 spark-shell 写代码,虽然方便,但有几个问题:

  • 代码写一行执行一行,不能保存,关掉就没了
  • 不能写复杂的多文件项目
  • 不能调试
  • 不能打包提交到集群运行

实际工作中,都是用 IDE(集成开发环境) 来写 Spark 程序的,最常用的就是 IntelliJ IDEA。

项目场景 ​

某竞赛网站经过几年的运营,保存了众多用户对该网站的访问日志数据。为了研究用户的兴趣爱好并改善用户体验,需要对网站用户的访问日志数据进行分析。

数据文件:raceData.csv(2023年5月至2024年2月的用户访问日志)

本项目要做的事:

  1. 在 IDEA 中搭建 Spark 开发环境
  2. 编写 Spark 程序统计每月访问量
  3. 用自定义分区器按年份分区保存结果
  4. 打包提交到集群运行

1.2 spark-shell vs IDE 编程 对比 ​

对比项spark-shellIDEA 编程
使用方式命令行交互式编写完整程序
代码保存不能保存保存为文件,可复用
SparkContext自动创建(sc变量)需要手动创建
适合场景学习、调试、小数据测试正式开发、生产环境
运行方式直接在shell里跑本地运行 / 打包提交集群
调试能力弱强(断点、单步等)
项目管理不支持支持多文件、多模块

💡 理解:spark-shell 就像"草稿纸",IDEA 就像"正式作业本"。学习用 shell,开发用 IDEA。


第二部分:搭建 Spark 开发环境(IDEA + Scala 插件) ​

2.1 安装 IntelliJ IDEA ​

下载 ​

  • 官网下载 IntelliJ IDEA 安装包(社区版即可,免费)
  • 版本:ideaIC-2022.3.3.exe

安装步骤 ​

  1. 双击安装包,进入安装向导
  2. 设置安装目录
  3. 完成一系列设置后,点击 Install
  4. 安装完成后重启电脑

首次启动 ​

  • 双击桌面图标启动
  • 第一次启动问是否导入以前的设定 → 选"不导入"
  • 欢迎界面可以切换主题(黑色/白色)

2.2 安装 Scala 插件 ​

IDEA 默认不支持 Scala,需要安装插件。有两种安装方式:

方式一:在线安装(推荐,有网时用) ​

  1. 打开 IDEA,欢迎界面左侧选 Plugins
  2. 在 Marketplace 搜索框输入 Scala
  3. 点击 Scala 右侧的 Install 按钮
  4. 安装完成后点击 Restart IDE 重启

方式二:离线安装(没网时用) ​

  1. 提前下载好 Scala 插件包(如 scala-intellij-bin-2022.3.3.zip)
  2. 在 Plugins 选项卡,点击齿轮图标 → 选 Install Plugin from Disk...
  3. 选择本地插件包路径
  4. 点击 OK 安装,然后重启 IDEA

⚠️ 注意:插件版本要和 IDEA 版本对应,否则可能装不上或用不了。


2.3 测试 Scala 插件 ​

创建第一个 Scala 项目 ​

  1. 欢迎界面选 New Project
  2. 填写项目信息:
    • Name:HelloWorld
    • Location:自定义存放路径
    • Language:Scala
    • Build System:IntelliJ
    • JDK:1.8.0_461
    • Scala SDK:2.13.16
  3. 点击 Create 创建

项目结构 ​

HelloWorld/
├── src/              # 源代码目录
├── .idea/            # IDEA配置文件(不用管)
└── HelloWorld.iml    # 项目配置文件(不用管)

编写 HelloWorld ​

  1. 右键 src → New → Package → 包名 com.tipdm.scalaDemo
  2. 右键包名 → New → Scala Class → 类名 HelloWorld,类型选 Object
  3. 编写代码:
scala
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 安装通义灵码插件 ​

  1. 打开 IDEA → Plugins → Marketplace
  2. 搜索 通义灵码(或英文 Lingma)
  3. 找到 "Lingma - Alibaba Cloud AI Coding Assistant"
  4. 点击 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 包,要先导入到项目中。

操作步骤 ​

  1. 打开 Project Structure(快捷键:Ctrl + Alt + Shift + S)
    • 或者菜单栏 File → Project Structure
  2. 选 Libraries 选项卡
  3. 点击 + 按钮 → 选 Java
  4. 找到本地 Spark 安装目录下的 jars 文件夹
  5. 点击 OK,Spark 依赖包就导入了

💡 说明:Spark 的 jars 文件夹里有所有需要的依赖包,一次性全部导入。


4.2 SparkContext 手动创建 ​

在 spark-shell 中,sc 变量是自动创建好的。 但在 IDEA 中写程序,需要手动创建 SparkContext。

创建步骤 ​

  1. 创建 SparkConf 对象:设置应用名称、运行模式等配置
  2. 创建 SparkContext 对象:传入 SparkConf 作为参数
  3. 用完后调用 sc.stop() 关闭

代码模板 ​

scala
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 ​

scala
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 程序,需要设置三个东西:

  1. 运行模式(必须设置,否则找不到 master)
  2. Hadoop bin 路径(Windows 需要)
  3. 自定义输入参数(如果程序需要参数)

5.1 设置运行模式 ​

方式一:在代码中设置(推荐,学习时用) ​

scala
val conf = new SparkConf().setAppName("WordCount").setMaster("local")

方式二:在 IDEA 运行配置中设置 ​

  1. 菜单栏 Run → Edit Configurations...
  2. 找到 VM options 输入框
  3. 输入:-Dspark.master=local
  4. 点击 OK

💡 两种方式的区别:

  • 代码中设置:写死了,打包到集群运行时要改代码
  • VM options 设置:不写死在代码里,更灵活
  • 实际开发:代码里不写 setMaster,打包后通过 spark-submit 指定

5.2 指定 Hadoop bin 路径(Windows 特有) ​

在 Windows 本地运行 Spark 程序,需要指定 Hadoop 的 bin 文件夹路径。

方式一:在代码中设置 ​

scala
System.setProperty("hadoop.home.dir", "E:\\hadoop-3.3.6")

方式二:设置 Windows 环境变量 ​

  • 在 Path 环境变量中添加 Hadoop 的 bin 文件夹路径
  • 这样程序运行时就不用每次指定了

还需要两个文件 ​

Hadoop 的 bin 文件夹里需要有:

  • winutils.exe
  • hadoop.dll

这两个文件可以从 GitHub 上下载对应版本的。

💡 说明:Linux 环境下不需要设置这个,只有 Windows 本地开发时才需要。


5.3 设置自定义输入参数 ​

如果程序需要从外部传入参数(比如输入路径、输出路径),就要设置 Program arguments。

代码中接收参数 ​

scala
val input = args(0)   // 第一个参数:输入路径
val output = args(1)  // 第二个参数:输出路径

IDEA 中设置参数 ​

  1. Run → Edit Configurations...
  2. 找到 Program arguments 输入框
  3. 输入参数(空格分隔),比如:E:\data\words.txt E:\data\output
  4. 点击 OK

5.4 本地运行完整示例 ​

代码(带参数版本) ​

scala
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 来指定。

scala
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 ​

  1. File → Project Structure(Ctrl+Alt+Shift+S)
  2. 选 Artifacts 选项卡
  3. 点击 + → JAR → Empty
  4. Name 填 JAR 包名称(如 word)
  5. 右侧找到你的项目 → 双击 '项目名' compile output
    • 它会移到左侧,表示已添加到 JAR 包中
  6. 点击 OK

步骤二:编译生成 JAR ​

  1. 菜单栏 Build → Build Artifacts...
  2. 选择你的 JAR 包名 → Build
  3. 等待编译完成

步骤三:找到 JAR 包 ​

  • 生成的 JAR 包在工程目录的 out/artifacts/ 目录下
  • 右键 JAR 包 → Open In → Explorer 可以打开文件位置

6.4 spark-submit 提交命令 ​

基本格式 ​

bash
spark-submit \
  --master <master-url> \
  --deploy-mode <deploy-mode> \
  --class <main-class> \
  <application-jar> \
  [application-arguments]

常用参数说明 ​

参数说明示例
--master集群地址(运行模式)spark://master:7077 / yarn
--deploy-mode部署模式:client / clustercluster
--class主类入口(包名+类名)com.tipdm.scalaDemo.WordCount
--conf任意 Spark 配置属性--conf spark.executor.memory=1g
--executor-memory每个 Executor 的内存1G
--executor-cores每个 Executor 的 CPU 核数1
--num-executorsExecutor 数量2
application-jarJAR 包路径/opt/data/word.jar
application-arguments传给 main 方法的参数输入路径 输出路径

示例1:Standalone 模式提交 ​

bash
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 模式提交 ​

bash
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:带资源配置的提交 ​

bash
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/wordcount2

6.5 配置优先级(面试常问) ​

Spark 配置可以在多个地方设置,优先级从高到低:

优先级配置位置说明
⭐⭐⭐⭐⭐ 最高代码中 conf.set()写死在代码里的
⭐⭐⭐⭐spark-submit 命令行参数--conf、--executor-memory 等
⭐⭐⭐配置文件 spark-defaults.confSpark 安装目录 conf 下
⭐⭐ 最低系统默认值Spark 内置默认值

💡 记忆口诀:代码 > 命令行 > 配置文件 > 默认值

越"近"的优先级越高,代码里写死的最高,命令行其次,配置文件再次,默认值最低。


6.6 提交到集群运行完整流程 ​

步骤 ​

  1. 启动集群:启动 Hadoop 和 Spark
  2. 上传 JAR 包:把 Windows 的 JAR 包传到 Linux 的 /opt/data/ 目录
  3. 上传数据:把数据文件上传到 HDFS
  4. 提交运行:用 spark-submit 提交
  5. 查看结果:在 HDFS 上查看输出

完整命令示例 ​

bash
# 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 的计算结果保存到内存/磁盘,下次用直接取,不用重算。

举个例子 ​

scala
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(存储级别)可指定存储级别,更灵活

代码示例 ​

scala
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(最近最少使用) 策略清除最老的分区
  • 内存中的缓存不是永久的,可能被清掉
  • 磁盘上的缓存不会被清

③ 手动清除缓存 ​

scala
// 手动清除不需要的缓存数据
data.unpersist()

💡 什么时候用持久化?

  • 一个 RDD 被多次使用时(比如迭代计算、机器学习)
  • 计算代价很高的 RDD(计算一次很费时间)
  • 数据量不大,能放进内存的

第八部分:Spark 数据分区 ​

8.1 为什么要控制数据分区? ​

Spark RDD 是由多个分区组成的,数据分布在多个节点上。 分布式计算中,网络通信的代价很大。 控制数据分区、减少网络传输(Shuffle),是提升性能的重要手段。

分区的两个维度 ​

  1. 分区个数:有多少个分区
  2. 分区方式:数据怎么分配到各个分区(分区器)

💡 理解:分区就像"分箱子装东西",分区个数是"有几个箱子",分区方式是"按什么规则往箱子里放"。


8.2 分区器(Partitioner) ​

重要前提 ​

只有键值对 RDD 才能设置分区方式(分区器) 非键值对 RDD 的分区方式是 None(只能设置分区个数)。

Spark 内置的两种分区器 ​

分区器说明特点
HashPartitioner哈希分区根据 Key 的哈希值取模分区,简单常用
RangePartitioner范围分区把一定范围的数据放到一个分区,排序后均匀分布

查看分区方式 ​

scala
// 查看RDD的分区器
rdd.partitioner   // 返回Option[Partitioner]

// 判断是否有分区器
rdd.partitioner.isDefined

// 获取具体的分区器
rdd.partitioner.get

8.3 partitionBy() 方法 ​

设置分区方式用 partitionBy() 方法,传入一个分区器。

scala
// 对键值对RDD设置哈希分区,3个分区
val partitioned = kv_rdd.partitionBy(new HashPartitioner(3))

// 查看分区数
partitioned.partitions.size

8.4 自定义分区器(重点+难点) ​

内置分区器满足不了需求时,可以自定义分区器。

实现步骤 ​

继承 org.apache.spark.Partitioner 类,实现 3 个方法:

方法签名作用
numPartitionsInt返回分区个数
getPartition(key: Any)Int根据 Key 返回分区 ID(0 到 numPartitions-1)
equals(other: Any)Boolean判断两个分区器是否相等

示例:按奇偶性分区 ​

自定义一个分区器,把偶数 Key 放分区 0,奇数 Key 放分区 1。

scala
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
  }
}

使用自定义分区器 ​

scala
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()
  }
}

提交运行 ​

bash
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() 方法 ​

语法:

scala
coalesce(numPartitions: Int, shuffle: Boolean = false)

参数说明:

  • numPartitions:重分区后的分区数
  • shuffle:是否进行 Shuffle,默认 false

规则:

  • shuffle = false(默认):只能减少分区数,不能增加(增加的话不生效)
  • shuffle = true:可以增加也可以减少分区数(会触发 Shuffle)

② repartition() 方法 ​

语法:

scala
repartition(numPartitions: Int)

本质:就是 coalesce(numPartitions, shuffle = true) 的简写。 特点:一定会触发 Shuffle,可以增加也可以减少分区数。


代码示例 ​

scala
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:计算竞赛网站每月的访问量 ​

需求 ​

统计每个月的网站访问量。

实现思路 ​

  1. 读取原始数据
  2. 取出时间字段(第5列,索引4)
  3. 统一日期格式,提取"年-月"
  4. 映射为 (年-月, 1) 的键值对
  5. 按 Key 求和 → 得到每月访问量

实现代码 ​

scala
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(_ + _)按月分组求和

提交运行 ​

bash
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)

实现思路 ​

  1. 先统计每月访问量(和任务4.1一样)
  2. 自定义分区器,按年份分区
  3. 用 partitionBy 应用分区器
  4. 保存结果到 HDFS

实现代码 ​

scala
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 写死了,如果年份更多,可以提取年份后判断:

scala
override def getPartition(key: Any): Int = {
  val year = key.toString.split("-")(0).toInt
  if (year == 2024) 1 else 0
}

提交运行 ​

bash
spark-submit \
  --master spark://master:7077 \
  --class race \
  /opt/data/race.jar \
  /tipdm/data/raceData.csv \
  /tipdm/data/race

查看结果 ​

bash
# 查看输出目录
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 排错通用思路 ​

  1. 看报错信息:先看异常类型和错误消息
  2. 本地先测通:先在本地模式(local)跑通,再打包提交集群
  3. 看 Spark UI:4040 端口看 Job、Stage、Executor 情况
  4. 看日志:YARN 日志或 Worker 日志找具体错误
  5. 小数据测试:先用小数据验证逻辑

第十一部分:实习 / 面试高频考点 ​

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 类 实现三个方法:

  1. numPartitions:返回分区个数
  2. getPartition(key):根据Key返回分区ID(0到numPartitions-1)
  3. 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? ​

  1. File → Project Structure → Artifacts → + → JAR → Empty
  2. 设置名称,添加 compile output
  3. Build → Build Artifacts → Build
  4. 去 out/artifacts 目录下找 JAR 包

Q15:Windows 本地运行 Spark 程序需要注意什么? ​

  1. 要设置 setMaster("local")
  2. 要设置 hadoop.home.dir
  3. Hadoop 的 bin 目录下要有 winutils.exe 和 hadoop.dll
  4. 只是本地测试用,生产环境代码里不要写死这些

Q16:怎么查看 Spark 程序运行情况? ​

  1. Spark Web UI:4040 端口(本地运行时)
  2. Standalone 集群:8080 端口看 Master UI
  3. YARN 集群:8088 端口看 ResourceManager UI
  4. 看日志: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):小区名称(用于过滤空值)

实现代码 ​

scala
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)
  }
}

提交运行 ​

bash
# 上传数据到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月

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