Spark3.x指北——2:Spark Core
本章为Spark Core的详细内容
Spark3.x指北全系列目录:
Spark基础概念请看:Spark3.x指北——1:Spark基础概念
SparkSQL内容请看:Spark3.x指北——3:SparkSQL
SparkStreaming内容请看:Spark3.x指北——4:SparkStreaming
目录
(1)我们首先需要一个服务器,用来等待客户端连接、接收数据、输出数据:
5.1.2 将客户端发送的数据变为发送一个计算任务(数据 & 操作硬编程)
(1)现在我们需要额外添加一个计算任务类,用于存储需要计算的数据 & 需要进行的计算操作:
(2)然后我们将客户端发送内容进行更改,通过ObjectOutputStream发送计算任务对象:
(3)最后,我们在服务器通过ObjectInputStream接收这个计算任务,然后执行这个任务里的计算方法,得出计算结果即可:
(2)保留5.1.2的提交任务类,这个类将会为Task提供数据和操作:
(3)我们在客户端将SubTask切分为两部分,分别赋值给Task1和Task2,然后将这两个任务发送给两台服务器9999和9998:
(4)服务器端分别设置两个端口接收客户端发送的两个任务,分别执行这两个任务然后输出结果:
(2)spark向Yarn申请资源,创建调度节点(Driver)和计算节点(Executor):
(3)Spark 框架根据需求将计算逻辑根据RDD分区划分成不同的任务:
(4)调度节点根据划分好的任务以及计算节点状态(是否就绪、节点距离等),分发给计算节点,完成计算:
Ⅰ 路径参数默认以当前环境根目录问基准,绝对路径 / 相对路径均可
Ⅳ 路径不仅可以为本地文件,同时也可以是HDFS、HBase中的数据
a 基础使用——将RDD数据集中每一个数据都执行 * 2 的转换
a 基础使用——一次性对一个分区的数据进行 * 2 的转换操作
c 两个算子:map & mapPartitions の对比
(5)glom算子——分区数据集类型的转换:List => Array
c 选出apache.log中 2015年5月17日的请求路径
(10)coalesce算子——改变分区数量(默认用于缩减分区)
(12)sortBy算子——根据规则函数对RDD数据进行排序
(13)针对两个RDDs数据的操作——交 & 并 & 差 & 拉链
d 区分:reduceByKey & groupByKey的区别和使用场景
e aggregateByKey算子——对分区内和分区间分别指定计算规则
f foldByKey算子——分区间和分区内计算规则相同的简化
g combineByKey算子——直接将第一个数据转换作为初始值进行计算
h 区分:reduceByKey & aggregateByKey & foldByKey & combineByKey的区别以及使用场景
i join & leftOuterJoin & rightOuterJoin算子——连接操作
② leftOuterJoin & rightOuterJoin算子
(2)collect算子——将RDD数据收集到内存形成结果数组
(6)takeOrdered算子——返回RDD中按照排序规则排序后的数据的前n个
(7)aggregate算子——分别与初始值进行分区内 & 分区间运算
(8)fold算子——aggregate算子的简化,分区内 & 分区间计算规则相同
(9)countByValue & countByKey算子——统计value出现的次数& 统计Key出现的次数
b saveAsObJectFile——将数据序列化为对象保存到文件中:
c saveAsSequenceFile——将KV类型转换为Sequence保存到文件中:
5.2.13 文件的读取 & 保存(保存具体源码参考5.2.8-save相关算子)
③ value方法——将当前Executor的最终累加结果进行返回
④ merge方法——将多个Executor的累加器在Driver端进行合并
c 版本三:对版本二的改进,使用ACC避免reduceByKey操作
5.4.3 需求二:Top10热门品类中每个品类的Top10活跃Session统计
5 Spark Core
Spark 计算框架为了能够进行高并发和高吞吐的数据处理,封装了三大数据结构,用于 处理不同的应用场景。
三大数据结构分别是:
- RDD : 弹性分布式数据集
- 累加器:分布式共享只写变量
- 广播变量:分布式共享只读变量
5.1 引入——通过网络编程模拟分布式计算
5.1.1 基础的客户端向服务端发送数据演示
(1)我们首先需要一个服务器,用来等待客户端连接、接收数据、输出数据:
//模拟服务器,用来执行任务
object Executor {
def main(args: Array[String]): Unit = {
//服务器配置
val server = new ServerSocket(9999)
//等待客户端链接
println("服务器启动,等待客户端链接")
val client: Socket = server.accept() //无程序连接就会阻塞
//接收客户端的数据
val in = client.getInputStream
val data = in.read()
println("接收到客户端数据:" + data)
in.close()
client.close()
server.close()
}
}
(2)然后我们需要一个客户端,向服务器发送数据:
//模拟驱动器,用于向服务器发送任务
object Driver {
def main(args: Array[String]): Unit = {
//连接服务器
val client = new Socket("localhost", 9999)
//向服务器发送信息
val out = client.getOutputStream
out.write(2)
out.flush()
out.close()
println("客户端发送完毕,关闭客户端")
client.close()
}
}
5.1.2 将客户端发送的数据变为发送一个计算任务(数据 & 操作硬编程)

(1)现在我们需要额外添加一个计算任务类,用于存储需要计算的数据 & 需要进行的计算操作:
class SubTask extends Serializable {
//需要计算的数据
val data: List[Int] = List(1, 2, 3, 4)
//计算所要进行的操作,用一个匿名函数来存储
val operation: (Int) => Int = _ * 2
//计算函数
def compute(): List[Int] = {
//将数据data全部进行operation操作,然后返回值作为结果
data.map(operation)
}
}
//注意,这个类必须混入Serializable特质,这样才能被对象输入流 & 输出流解析!
(2)然后我们将客户端发送内容进行更改,通过ObjectOutputStream发送计算任务对象:
//模拟驱动器,用于向服务器发送任务
object Driver {
def main(args: Array[String]): Unit = {
//连接服务器
val client = new Socket("localhost", 9999)
//向服务器发送对象
val out = client.getOutputStream
val objOut = new ObjectOutputStream(out)
val subTask = new SubTask
objOut.writeObject(subTask)
objOut.flush()
objOut.close()
println("客户端发送完毕,关闭客户端")
client.close()
}
}
(3)最后,我们在服务器通过ObjectInputStream接收这个计算任务,然后执行这个任务里的计算方法,得出计算结果即可:
//模拟服务器,用来执行任务
object Executor {
def main(args: Array[String]): Unit = {
//服务器配置
val server = new ServerSocket(9999)
//等待客户端链接
println("服务器启动,等待客户端链接")
val client: Socket = server.accept() //无程序连接就会阻塞
//接收客户端的任务对象并执行
val in = client.getInputStream
val objIn = new ObjectInputStream(in)
val task = objIn.readObject().asInstanceOf[SubTask]
println("开始执行计算任务")
val results = task.compute()
println("接收到客户端数据:" + results)
objIn.close()
client.close()
server.close()
}
}
//我们需要通过asInstanceOf将接收的对象强转为SubTask
5.1.3 将计算任务进行分布式处理
观察我们在 5.1.2 实现的代码,是将一个完整的计算任务类发送给了服务器,让其执行。这显然不是一个分布式计算的逻辑。我们需要实现一个真正的分布式计算,也就是将这个完整的任务进行拆分,分发给多个任务进行执行。

(1)我们重新定义一个任务类,这个任务类的 计算数据 和 计算操作 均不赋值,只是封装一个compute方法,这个方法中对 数据 进行 操作。我们可以将这个类理解为控制抽象,仅提供对数据的规范和对操作的规范,具体的数据和操作均由调用者提供:
class Task extends Serializable {
var data: List[Int] = _
var operation: (Int) => Int = _
def compute(): List[Int] = data.map(operation)
}
(2)保留5.1.2的提交任务类,这个类将会为Task提供数据和操作:
class SubTask extends Serializable {
//需要计算的数据
val data: List[Int] = List(1, 2, 3, 4)
//计算所要进行的操作,用一个匿名函数来存储
val operation: (Int) => Int = _ * 2
//计算函数
def compute(): List[Int] = {
//将数据data全部进行operation操作,然后返回值作为结果
data.map(operation)
}
}
(3)我们在客户端将SubTask切分为两部分,分别赋值给Task1和Task2,然后将这两个任务发送给两台服务器9999和9998:
//模拟驱动器,用于向服务器发送任务
object Driver {
def main(args: Array[String]): Unit = {
//连接服务器
val client1 = new Socket("localhost", 9999)
val client2 = new Socket("localhost", 9998)
//提供完整的需要进行的任务
val subTask = new SubTask
//向服务器发送对象
//将计算任务拆分为两部分,分别发送给两个服务器
//任务1,计算subTask.data的前两个元素,并且发送给client1
val objOut1 = new ObjectOutputStream(client1.getOutputStream)
val task1 = new Task
task1.data = subTask.data.take(2)
task1.operation = subTask.operation
objOut1.writeObject(task1)
objOut1.flush()
objOut1.close()
client1.close()
//任务2,计算subTask.data的后两个元素,并且发送给client2
val objOut2 = new ObjectOutputStream(client2.getOutputStream)
val task2 = new Task
task2.data = subTask.data.takeRight(2)
task2.operation = subTask.operation
objOut2.writeObject(task2)
objOut2.flush()
objOut2.close()
client2.close()
println("客户端发送完毕,关闭客户端")
}
}
(4)服务器端分别设置两个端口接收客户端发送的两个任务,分别执行这两个任务然后输出结果:
//模拟服务器,用来执行任务
object Executor {
def main(args: Array[String]): Unit = {
//服务器配置,两个服务器,分别接收两个客户端的任务
val server1 = new ServerSocket(9999)
val server2 = new ServerSocket(9998)
//等待客户端链接
println("服务器启动,等待客户端链接")
val client1: Socket = server1.accept() //无程序连接就会阻塞
val client2: Socket = server2.accept()
//接收客户端的任务对象并执行
//执行第一个任务
val in1 = client1.getInputStream
val objIn1 = new ObjectInputStream(in1)
val task1 = objIn1.readObject().asInstanceOf[Task]
println("开始执行计算任务1")
val results1 = task1.compute()
println("计算结果1:" + results1)
objIn1.close()
client1.close()
server1.close()
println("服务器1关闭")
//执行第二个任务
val objIn2 = new ObjectInputStream(client2.getInputStream)
val task2 = objIn2.readObject().asInstanceOf[Task]
println("开始执行计算任务2")
val result2 = task2.compute()
println("计算机给2:" + result2)
objIn2.close()
client2.close()
server2.close()
println("服务器2关闭")
}
}
(5)对于5.1.3分布式计算模拟的总结
- 我们需要一个任务类Task,这个类中仅仅存放对 计算数据的规范(比如需要进行计算的数据是List[Int]) 以及 计算操作的规范(比如这个操作是一个(Int) => Int 的匿名函数,用于通过map方法进行调用),这可以理解为一个抽象控制,等待具体任务传递数据。
- 我们需要一个有具体的任务SubTask,这个任务有任务类所需要的具体数据以及具体操作,然后我们在客户端中将这个具体任务切分为多个部分,赋值给多个Task,将这些Task分发给不同的服务器。
- 在服务器中,我们接收到了不同客户端发来的不同Task,我们可以对这些Task执行计算操作,得到各个部分的结果。最后,我们可以考虑将这些结果整合到一起,或者分别输出。
5.2 数据结构——RDD
5.2.1 什么是RDD
RDD(Resilient Distributed Dataset)叫做弹性分布式数据集,是 Spark 中最基本的数据 处理模型。代码中是一个抽象类,它代表一个弹性的、不可变、可分区、里面的元素可并行 计算的集合。
我们可以将RDD与5.1.3中的SubTask进行类比,SubTask中封装了计算逻辑和准备进行计算的数据,通过封装的计算逻辑分发给位于不同服务器上的节点,进行分布式的计算。我们对于集合取左和取右元素的操作就类似于对其进行分区。

- 弹性
- 存储的弹性:内存与磁盘的自动切换;
- 容错的弹性:RDD仅记录计算逻辑,不保存数据,因此数据丢失可以自动恢复(读取计算逻辑,重新进行计算);
- 计算的弹性:计算出错重试机制;
- 分片的弹性:可根据需要重新分片。
- 分布式:数据存储在大数据集群不同节点上,毕竟spark就是基于分布式数据存储进行分布式计算的。
- 数据集:RDD封装了计算逻辑,并不保存数据。RDD通过保存数据路径来获取所需要的数据,比如内存数据的引用、父类RDD的引用等等。
- 数据抽象:RDD是一个抽象类,需要子类具体实现
- 不可变:RDD封装了计算逻辑,是不可以改变的,想要改变,只能产生新的RDD,在 新的RDD里面封装计算逻辑。这是RDD需要通过装饰者模式进行功能添加的基础。
- 可分区、并行计算:一个RDD内可以被分为多个区,对一个完整数据进行拆分计算。
5.2.2 从IO的装饰者模式到RDD的组合
我们知道,IO流被分为许多类,有最基础的字节流Stream,还有带有缓冲区的BufferStream,以及字符流、对象流等等。
这些流的核心,都是最基础的字节流。比如BufferStream,是在字节流读取的基础上,加上了一个缓冲区,等缓冲区读满再输出;又或者字符缓冲区流,在字节流读取字节的基础上转换为字符,然后将字符放到缓冲区,缓冲区读满输出…
这些流在创建对象的过程都有一个共同特点,就是在原有对象上添加新的对象,让原有对象拥有新的功能,而不改变原有的属性。比如BufferStream,创建的过程类似于:new BufferStream(new InputStream),相当于为InputStream流添加了BufferStream的装饰,使得这个InputStream流可以通过缓冲区进行读取:

这种设计模式就被称作装饰者模式,通过将一些需要的对象作为逻辑封装到原有对象上,以实现对原有对象的拓展,但不会改变原有对象的属性。这与scala中的隐式类的创建逻辑上是类似的,原有对象并不会被改变类型,但是同时也能够具有新的特性。
这与RDD是有异曲同工的。因为RDD是不可变的,所以我们想要在上一个RDD的基础上进行进一步处理,我们就需要通过装饰者模式为其添加新的功能。比如我们想要一个进行map操作的RDD,那么这个RDD就会在基础的读取文件RDD对象上添加mapRDD的装饰。这个mapRDD会在读取文件RDD结束读取后的基础上添加map操作,成为一个可以进行map的RDD。
需要注意,对RDD添加map、flat等转换操作(即生成另外的RDD的操作)时,这些方法并不会立即执行,只有当我们调用collect、saveAsTextFile等行动操作(生成我们需要的数据类型)时,才会将之前的转换操作装饰执行。我们可以将进行转换操作看作规划导航路线,但我们此时仅仅是规划,只有当我们进行行动操作时,才是真正上路。

5.2.3 RDD的五大特性

(1)分区列表

一个RDD中可以划分为多个分区,将数据分别放到不同的分区进行分布式计算。这个分区列表是spark实现分布式计算的基础。
这个分区可以是系统指定,也可以我们直接指定。比如一个List(1, 3, 5, 7, 9, 11, 13, 15),我们可以指定2个元素为一个分区,将这个RDD中的List放到四个分区进行计算。
(2)分区计算函数

需要注意,RDD是最小的计算单元,也就是说RDD中的每一个分区执行的计算操作都是一样的,只是每个分区的计算数据不同。
(3)RDD依赖关系的获取(血统)

这是我们获取RDD之间依赖关系的重要函数,通过知道RDD之间的依赖,spark就可以更好根据依赖关系构成DAG,优化计算。同时,当中间RDD计算出错,也可以根据依赖关系,将RDD进行重新计算。
(4)分区器(KV数据类型时采用的分区策略)

(5)计算节点分发的首选位置

当一个计算节点与数据节点的位置最近,那么他们之间进行任务发送所需要的IO是最少的。如果一个数据源和RDD在同一个节点上,那么我们肯定是选择将RDD分区任务分发到这个节点上,直接进行计算,毕竟这样不会造成数据的移动,效率肯定是更高的。
5.2.4 RDD的执行原理
从计算的角度来讲,数据处理过程中需要计算资源(内存 & CPU)和计算模型(逻辑)。 执行时,需要将计算资源和计算模型进行协调和整合。
Spark 框架在执行时,先申请资源,然后将应用程序的数据处理逻辑分解成一个一个的 计算任务。然后将任务发到已经分配资源的计算节点上, 按照指定的计算模型进行数据计 算。最后得到计算结果。
RDD是Spark框架中用于数据处理的核心模型,接下来我们看看,在Yarn环境中,RDD 的工作原理:
(1)启动Yarn集群,获取资源调度框架,用于管理资源:

(2)spark向Yarn申请资源,创建调度节点(Driver)和计算节点(Executor):

(3)Spark 框架根据需求将计算逻辑根据RDD分区划分成不同的任务:

(4)调度节点根据划分好的任务以及计算节点状态(是否就绪、节点距离等),分发给计算节点,完成计算:

5.2.5 RDD基础编程——创建RDD
(1)从内存中直接创建RDD——以内存中数据为数据源
a 通过 parallelize 方法
object rdd_from_memory {
def main(args: Array[String]): Unit = {
//TODO 准备Spark环境
//准备配置参数
val sparkConf = new SparkConf()
.setMaster("local[*]") //使用本机的CPU核数
.setAppName("RDD_BUILD") //SparkJob的名称
//创建上下文环境对象
val sc = new SparkContext(sparkConf)
//TODO 创建RDD
//从内存中创建RDD,将内存中集合数据作为处理数据源
val rdd: RDD[Int] = sc.parallelize(List(1, 3, 5, 7, 9))
//执行collect方法,调用rdd执行
rdd.collect().foreach(println)
//TODO 关闭环境
sc.stop()
}
}
b 通过 makeRDD 方法创建(常用)
object rdd_from_memory {
def main(args: Array[String]): Unit = {
//TODO 准备Spark环境
//准备配置参数
val sparkConf = new SparkConf()
.setMaster("local[*]") //使用本机的CPU核数
.setAppName("RDD_BUILD") //SparkJob的名称
//创建上下文环境对象
val sc = new SparkContext(sparkConf)
//TODO 创建RDD
//从内存中创建RDD,将内存中集合数据作为处理数据源
//val rdd: RDD[Int] = sc.parallelize(List(1, 3, 5, 7, 9))
val rdd: RDD[Int] = sc.makeRDD(List(1, 3, 5, 7, 9))
//执行collect方法,调用rdd执行
rdd.collect().foreach(println)
//TODO 关闭环境
sc.stop()
}
}
实际上,makeRDD底层是直接调用了parallelize方法:

(2)从外部存储中创建RDD
a 通过textFile方法
object rdd_from_file {
def main(args: Array[String]): Unit = {
// TODO 准备spark环境
//准备配置
val sparkConf = new SparkConf().setMaster("local[*]").setAppName("rdd_from_file")
//准备上下文
val sc = new SparkContext(sparkConf)
// TODO 创建RDD——将文件中数据作为数据源
val rdd = sc.textFile("D:\\develop\\dataScience_learn\\spark_learn\\IDEA_PROJECT\\Spark_IDEA\\datas\\1.txt") //路径默认以当前环境根目录为基准,可以写绝对路径 / 相对路径
rdd.collect().foreach(println)
// TODO 关闭环境
sc.stop()
}
}
Ⅰ 路径参数默认以当前环境根目录问基准,绝对路径 / 相对路径均可
比如我的项目结构如下:

我想要读取datas中1.txt的文件内容,那么通过绝对路径读取为:
sc.textFile("D:\\develop\\dataScience_learn\\spark_learn\\IDEA_PROJECT\\Spark_IDEA\\datas\\1.txt")
通过相对路径读取为:
val rdd = sc.textFile("datas/1.txt")
Ⅱ 当我只写目录名时,我就可以读取目录中所有文件
仍然以a中项目结构为准,此时我想要读取datas中所有文件内容,我就可以将textFile参数写作:
sc.textFile(“datas”)
Ⅲ 可以使用通配符,筛选目录中的文件
当我想要读取datas目录中所有以1开头的文件,我就可以通过”datas/1*.txt”对要读取的文件进行筛选过滤:
val rdd = sc.textFile("datas/1*.txt")
Ⅳ 路径不仅可以为本地文件,同时也可以是HDFS、HBase中的数据
sc.textFile(“hdfs://node1:8020/datas”)
b 通过wholeTextFiles方法创建
textFile方法是以行为单位读取,并不关心数据来源。当我们试图读取一个目录中的所有文件,其输出通常为:

若我们想要知道数据来自哪个文件,并且可以在控制台进行区分,我们就可以使用wholeTextFiles方法:
val rdd = sc.wholeTextFiles("datas")
输出为:
(file:/D:/develop/dataScience_learn/spark_learn/IDEA_PROJECT/Spark_IDEA/datas/1.txt,Hello World
Hello Spark
darren is learning spark)
(file:/D:/develop/dataScience_learn/spark_learn/IDEA_PROJECT/Spark_IDEA/datas/2.txt,Hello World
Hello Spark
darren is learning spark)
(file:/D:/develop/dataScience_learn/spark_learn/IDEA_PROJECT/Spark_IDEA/datas/11.txt,)
(3) 从其他RDD创建
主要是通过一个RDD运算完后,再产生新的RDD。详情请参考后续章节
(4) 直接创建RDD(new)
使用new的方式直接构造RDD,一般由Spark框架自身使用。
5.2.6 RDD基础编程——分区与并行度
我们在RDD基础概念中已经提到,RDD是可以具有分区的,也就是将数据划分到不同分区中进行分布式计算。因此我们在创建RDD时,也是可以主动指定RDD的分区的。
(1)从内存中读取数据时创建分区
a 分区的设定
object rdd_parallelize {
def main(args: Array[String]): Unit = {
// TODO 准备环境
val sparkConf = new SparkConf().setMaster("local[*]").setAppName("RDD")
val sc = new SparkContext(sparkConf)
// TODO 创建RDD
// RDD并行度 & 分区
val rdd = sc.makeRDD(
List(1, 2, 3, 4), //数据源
2 //分区数量
)
rdd.saveAsTextFile("output") //按照设定分区保存到分区文件中
// TODO 关闭环境
sc.stop()
}
}
若我们不指定分区,这个分区参数则具有默认值,这个默认值是环境可以使用的最大CPU核数:

我们需要注意,分区与并行度并不能画上等号。分区是我们要将任务划分的区域数,并行度是同一时间可以同时执行的任务数。比如我们有8个分区,但是当前可用核数为4,那么当前仅有4个任务可以并行,并行度为4。
b 分区数据的分配
我们需要对集合数据如何分区进行一个讨论。对于可以被平均分配分区数据我们先不讨论,原理是共同的,我们通过一个不可被平均分配分区的数据来查看RDD对于内存数据分区的策略。
假设我们有一个List[5],而makeRDD指定的分区数为3,这显然是一个不可平均分配的分区:
//进行RDD创建——基于内存
val rdd = sc.makeRDD(List(1, 2, 3, 4, 5), 3)
//行动操作,将数据保存到分区文件
rdd.saveAsTextFile("output")
对于这样的分区操作,最终结果为:[1]、[2, 3]、[4, 5]。我们直接展示内存数据分区的核心——分区策略计算源码:

这是一个迭代的过程,每一次都返回当前分区所包含的数据范围,然后进行下一次数据范围迭代。数据进行迭代的范围为 0 until numSlices,也就是(0, 1, …, numSlices – 1)。注意:这里返回的数据范围是指集合数据的索引范围,并不是指数据本身!
比如我们以List(1, 2, 3, 4, 5)为例,指定numSlices = 3。那么迭代范围即为 0 until 3,迭代的数据分别为 0, 1, 2。我们开始迭代:
- 第一次迭代:i = 0,start = (i * length) / numSlices = (0 * 5) / 3 = 0,end = ((i + 1) * length) / numSlices = ((0 + 1) * 5) / 3 = 1,因此第一个范围为(0, 1)。
- 第二次迭代:i = 1,则start = (1 * 5) / 3 = 1,end = ((1 + 1) * 5) / 3 = 3,第二个范围为(1, 3)。
- 第三次迭代:i = 2,则start = 3,end = 5,第三个范围为(3, 5) 。
我们在这三次迭代中得到的范围分别为(0, 1),(1, 3),(3, 5)。我们查看最后将数组进行切分的源码:

可以看到,这是以我们得到的分区范围,进行start until end 的切片。因此(0, 1)中,包含0索引元素,即List中的1;(1, 3)中,包含1, 2索引元素,即List中的2, 3;(3, 5)中,包含3, 4索引元素,即List中的4, 5。因此最终分区就是[1], [2, 3], [4, 5]。
(2)从文件中读取数据时创建分区
object rdd_parallelize_fromFile {
def main(args: Array[String]): Unit = {
// TODO 准备环境
val sparkConf = new SparkConf().setMaster("local[*]").setAppName("RDD")
val sc = new SparkContext(sparkConf)
// TODO 创建RDD
// RDD并行度 & 分区
val rdd = sc.textFile("datas/1.txt")
rdd.saveAsTextFile("output") //按照设定分区保存到分区文件中
// TODO 关闭环境
sc.stop()
}
}
同样的,当我们不指定分区数量,那么将采用默认的分区数(但这个值与makeRDD的默认值并不相同,需要区分!):


需要知道,这个默认分区是与文件大小强相关的!
5.2.7 RDD基础编程——转换算子

(1)map算子——对分区内数据进行依次处理 & 转换

将处理的数据逐条进行映射转换,这里的转换可以是类型的转换,也可以是值的转换。
map算子是对RDD中分区内每一个数据都进行转换操作,此操作会对每一个元素都进行转换,因此不会导致新RDD该分区数据量的改变。
a 基础使用——将RDD数据集中每一个数据都执行 * 2 的转换
object rdd_map_transform {
def main(args: Array[String]): Unit = {
val sparkConf = new SparkConf().setMaster("local[*]").setAppName("Operator")
val sc = new SparkContext(sparkConf)
// TODO map算子
val rdd = sc.makeRDD(List(1, 2, 3, 4))
val mapRDD = rdd.map(_ * 2)
val result = mapRDD.collect()
println(result.mkString(","))
sc.stop()
}
}
正如scala中的map是将 集合数据 通过 匿名转换方法 转换为一个新的集合数据,spark中的map算子也是将 原有的RDD数据 逐条转换为 新的RDD ,使用方法与scala是极其相似的。但需要注意,使用map方法仅仅返回一个RDD,需要调用行动算子collect等才能构成一个集合数据类型。
b 一个小实操案例
案例要求:
![]()
日志位于/datas下,且内容如下:

实现思路为——观察日志信息发现,每一行中都是”端口 - - 时间戳 +0000 请求方式 请求路径”的格式,也就是说每一个信息都是通过” “进行分割的。通过文件读取获取RDD时,是按行进行读取的,因此我们可以对每一行的字符串都通过split(“ ”),然后获取分割数组的最后一个元素即可:
object rdd_map_practice {
def main(args: Array[String]): Unit = {
//准备spark环境
val sparkConf = new SparkConf().setMaster("local[*]").setAppName("MAP")
val sc = new SparkContext(sparkConf)
//读取apache.log文件,用" "截取,并且去除最后一段字符串,即为我们需要的路径数据
val fileRDD = sc.textFile("datas/apache.log")
val result = fileRDD.map(
line => { //line为每一行的数据
val datas = line.split(" ") //每一行切分后的字符数组
datas(datas.length - 1) //将字符数组最后一个元素取出即为目标路径
}
).collect()
result.foreach(println)
sc.stop()
}
}
输出结果如下:

c 并行计算效果演示
RDD核心是将计算逻辑封装,然后分发后进行并行计算。对于不同分区的数据,采用并行计算;而对于同一分区,采用串行计算。
- 我们先从分区数为1的RDD入手:
val rdd = sc.makeRDD(List(1, 2, 3, 4), 1)
val mapRDD = rdd.map(
num => {
println("mapRDD>: " + num)
num
}
)
val mapRDD1 = mapRDD.map(
num => {
println("mapRDD1#: " + num)
num
}
)
mapRDD1.collect()
我们可能会认为,应该是mapRDD的逻辑先执行完,再执行mapRDD1的逻辑,但实际上该代码的输出结果为:

对于makeRDD(List(1, 2, 3, 4), 1),由于分区数为1,因此所有数据都位于一个分区中。在这种情况下,分区中的数据要一个一个执行计算逻辑,只有前面一个数据计算逻辑执行完毕,才会执行下一个数据的计算。也就是说,需要先对1进行mapRDD和mapRDD1的所有操作后,才会进行2的操作,然后是3,最后是4。
即:同一分区内,数据执行是有序的。
- 接着我们看分区数不为1的情况:
val rdd = sc.makeRDD(List(1, 2, 3, 4), 2)
val mapRDD = rdd.map(
num => {
println("mapRDD>: " + num)
num
}
)
val mapRDD1 = mapRDD.map(
num => {
println("mapRDD1#: " + num)
num
}
)
mapRDD1.collect()
此时,(1, 2)位于同一个分区,(3, 4)位于另一个分区。由于同一个分区内数据串行,所以mapRDD和mapRDD1都会先将1 / 3的计算执行,然后再执行2 / 4。但是对于两个分区之间,由于是并行计算,所以先执行对1的操作还是对2的操作,我们是无法确定的。输出结果如下:

即:不同分区之间,数据执行是无序的。
(2)mapPartitions算子——对分区进行整体处理

map方法为我们提供单个处理分区内数据的方法,但是由于map方法只能对分区数据进行单独处理(即分区内一个数据全部被处理完成后才能进行下一个数据处理),因此效率不高。spark为我们提供了mapPartitions方法,使得我们能够批量处理分区数据。
通过mapPartitions方法,将待处理的数据以分区为单位发送到计算节点进行处理,这里的处理是指可以进行任意的处 理,哪怕是过滤数据。
需要注意,这个方法是将一个 Iterator对象(即当前分区数据集) 转换为另一个 Iterator对象,因此在匿名函数返回中,一定要构建一个迭代器返回值。
Iterator => Iterator的要求是为了符合RDD的设计思想,即保存数据操作逻辑而非数据本身。Iterator => Iterator的操作实际上也是一个计算逻辑,在mapPartitions操作后,将会创建一个具有该操作的新RDD。
同时,由于返回值仅需要是一个Iterator,因此分区内元素的个数并不需要保证与原RDD相同,这是与map操作的另一个区别。
a 基础使用——一次性对一个分区的数据进行 * 2 的转换操作
val rdd = sc.makeRDD(List(1, 2, 3, 4), 2)
val mpRDD = rdd.mapPartitions(
iter => { //迭代器,是整个分区的数据
println(">>>>>>>>")
iter.map(_ * 2) //对迭代器中每个元素进行 * 2 的迭代操作,生成新迭代器并返回
}
).collect().foreach(println)
这是上述代码的输出结果:

若使用的是map方法,由于makeRDD指定了分区数为2,因此数据处理过程应该是 先处理两个分区中的第一个数据 1 / 3,输出 2 / 6,最后处理 2 / 4,输出4 / 8。但是显然,mapPartitions方法的输出是按分区进行输出,先处理了分区1的(1, 2),再处理分区2的(3, 4),这显然是将整个分区数据全部加载后再进行的处理。
但需要注意,虽然mapPartitions为我们提供了批量处理数据的方法,这个方法会将分区数据引用全部加载到内存,并且不会自动释放(JVM的垃圾回收懒处理机制)。若我们不进行主动内存释放,当数据量很大的情况下,内存很可能会溢出。
b 案例演示
案例需求:

代码实现:
val rdd = sc.makeRDD(List(1, 3, 5, 2, 1, 3, 10, 11, 8), 3)
rdd.mapPartitions(
iter => {
//获取每个分区数据的最大值
List(iter.max).iterator
}
).collect().foreach(println)
结果输出:

需要注意,Iterator的特性是仅能够使用一次,因此iter.max结束后,iter这个迭代器将不再具有元素。若想要在这个方法中多次使用iter,则需要提前将iter物化为集合,然后再进行处理。
c 两个算子:map & mapPartitions の对比
- 数据处理角度
Map 算子是分区内一个数据一个数据的执行,类似于串行操作。而mapPartitions算子 是以分区为单位进行批处理操作。
- 功能的角度
Map 算子主要目的将数据源中的数据进行转换和改变。但是不会减少或增多数据。 MapPartitions 算子需要传递一个迭代器,返回一个迭代器,没有要求的元素的个数保持不变, 所以可以增加或减少数据
- 性能的角度
Map 算子因为类似于串行操作,所以性能比较低,而是mapPartitions算子类似于批处 理,所以性能较高。但是mapPartitions算子会长时间占用内存,那么这样会导致内存可能 不够用,出现内存溢出的错误。所以在内存有限的情况下,不推荐使用。使用map操作。
(3)mapPartitionsWithIndex算子

这个方法会自动获取 调用者RDD的分区数据集合 以及 当前分区数据集所在的索引 ,因此对于其中匿名方法参数解析:(Int, Iterator[T]),第一个参数即为分区索引,第二个参数为当前分区所有数据的迭代器。
同样的,该方法的返回值也需要是一个迭代器。
a 通过该方法获取第二个分区的所有数据
//选取数据集中第二个分区的数据
val rdd = sc.makeRDD(List(1, 2, 3, 4, 5, 6), 3)
rdd.mapPartitionsWithIndex(
{
//对于第二个分区,返回数据
case (1, datas) => datas
//对于其他分区,返回空迭代器
case (_, datas) => Nil.iterator
}
).collect().foreach(println)
若使用偏函数进行模式匹配,则需要将所有可能情况都进行匹配。此处需要匹配的主要是索引,因此考虑索引为1或其余可能即可。需要注意,返回结果需要是一个迭代器,所以即使我们在其余情况返回一个Nil集合,也需要将其迭代器化。
该方法最终返回结果为分区索引为1的迭代器,遍历迭代器输出结果为:

b 通过该方法获取每个数据所在的分区索引
//获取对应数据所在的分区
val rdd = sc.makeRDD(List(1, 2, 3, 4, 5, 6))
rdd.mapPartitionsWithIndex(
{
case (index, datas) => datas.map((index, _))
}
).collect().foreach(println)
输出结果为:

(4)flatMap算子——扁平化

将处理的数据进行扁平化后再进行映射处理,所以算子也称之为扁平映射。
a 将RDD中两个数据集合进行扁平化合并
val rdd: RDD[List[Int]] = sc.makeRDD(List(List(1, 2, 3), List(4, 5)))
//将rdd中两个集合数据扁平化
val flatRDD: RDD[Int] = rdd.flatMap(
//形参为当前需要进行扁平化的集合
//返回值为封装扁平化后元素的集合
list => list
)
flatRDD.collect().foreach(println)
需要注意,flatMap中匿名函数形参和返回值都为list,但这两个list含义并不相同!形参中的list代表当前需要进行扁平化的集合,而返回值的list表示封装扁平化后元素的集合!
b 将RDD中的字符串集合元素按空格拆分并且扁平化
val rdd = sc.makeRDD(List("hello world", "hello scala", "hello spark"))
//将rdd中的单词按空格拆分并且扁平化为集合
val result = rdd.flatMap(_.split(" ")).collect()
println(result.mkString(","))
c 对于a的进阶,将混合类型的数据集合扁平化
现在RDD中有一个List(List, Int)类型的混合集合,我想要将其中所有元素都扁平化,放到一个集合中。
val rdd = sc.makeRDD(List(List(1, 2, 3), List(4, 5), 7, 8))
//将rdd中所有元素(集合类型 / 单独元素)合并到一起
val results = rdd.flatMap(
{
case list: List[Int] => list
case data => List(data)
}
).collect()
println(results.mkString(","))
直接在flatMap中通过偏函数进行类型模式匹配,当匹配到List[Int]类型数据,直接将其作为封装结果返回;当匹配到单独的元素数据,则将该元素封装为List返回。
(5)glom算子——分区数据集类型的转换:List => Array

将同一个分区的数据直接转换为相同类型的内存数组进行处理,分区不变。
a 基础使用
val rdd = sc.makeRDD(List(1, 2, 3, 4), 2)
//List => Int => Array
val glomRDD = rdd.glom()
glomRDD.collect().foreach(arr => println(arr.mkString(",")))
需要将该方法与flatMap进行区分,glom方法仅将每个分区数据都转换成数组,但是不会将分区间数据进行合并。
b 将分区内的最大值取出,进行分区间最大值求和
我们采用glom方法,先将分区数据转换为内存数组,然后通过操作内存数组求最值 & 求和:
val rdd = sc.makeRDD(List(1, 2, 3, 4), 2)
//选取分区内最大值,并且对分区间最大值进行求和
val result = rdd.glom().map(_.max).sum()
println(result)
当然还有另一种写法,通过mapPartitions获取分区整体数据,对整体求最值然后求和:
val result2 = rdd.mapPartitions(list => List(list.max).iterator).sum()
println(result2)
c 分区不变的含义
对于我们上述算子(map、mapPartitions…),都是在原本的分区进行数据操作。新RDD与旧RDD的数据实际上是通过pipeline进行传输的:

由于数据操作并没有跨区进行交互,只需要本分区的数据进行计算,因此这样的算子并不会造成数据分区的改变,这也被称作 窄依赖 。
(6)groupBy算子——分布式分组

a 基础使用
val rdd = sc.makeRDD(List(1, 2, 3, 4), 2)
//按奇偶进行分组
val result = rdd.groupBy(_ % 2).collect()
println(result.mkString(","))
b shuffle概念了解
将数据根据指定的规则进行分组, 分区默认不变,但是数据会被打乱重新组合,我们将这样的操作称之为shuffle。
极限情况下,数据可能被分在同一个分区中,一个组的数据在一个分区中,但是并不是说一个分区中只有一个组:

c 统计apache.log中每个时间段访问数据量
apache.log内容如下:

我们仅统计每天每个小时的数据访问量。实现如下:
//从apache.log中获取每个时间段的访问量
val rdd = sc.textFile("datas/apache.log")
//按时间分组,形成time -> 访问数据总量 的元组
val groupRDD = rdd.groupBy(line => line.split(" ")(3).split(":")(0) + ":" + line.split(" ")(3).split(":")(1))
//对分组结果进行数据量的统计
val result = groupRDD.map(kv => kv._1 -> kv._2.size).collect()
result.foreach(println)
注意,想要统计迭代器中元素的数量,可以直接通过iter.size获取。
(7)filter算子——过滤

传给该算子一个规则,该算子就会将RDD每个分区内符合规则的数据留下,不符合规则的数据去除,形成新的RDD。
a 基础使用
val rdd = sc.makeRDD(List(1, 2, 3, 4))
//将奇数元素过滤并留在rdd中
val results = rdd.filter(_ % 2 == 1).collect()
println(results.mkString(","))
b 数据倾斜
由于filter是对分区内数据进行筛选,筛选后并不会改变分区,因此可能会出现:某个分区数据大部分都符合规则,而另一个分区数据基本不符合规则,这就导致了整个RDD中某些分区数据量过大,而某些分区基本没有数据。这就是数据倾斜。
c 选出apache.log中 2015年5月17日的请求路径
//选出apache.log中 2015年5月17日的请求路径
val rdd = sc.textFile("datas/apache.log")
val results = rdd.filter(_.split(" ")(3).split(":")(0) == "17/05/2015").collect()
results.foreach(println)
(8)sample算子——抽样

参数说明:
withReplacement:抽取后是否放回。true——放回,false——不放回
fraction:每个元素期望被抽取的概率
seed:随机数种子
注意:
seed一旦确定,那么就随机数确定,因此每次抽取都是相同元素。若我们不指定seed,则将会采用当前系统时间作为seed。
使用场景:
当我们在shuffle打乱分区数据位置后,可能会出现数据倾斜的情况。这时我们就可以通过sample抽样,将分区内部分数据移动到数据量较小的分区中。
val rdd = sc.makeRDD(List(1, 2, 3, 4, 5, 6, 7, 8, 9, 10))
val results = rdd.sample(false, 0.5).collect()
println(results.mkString(","))
(9)distinct算子——分布式去重


a 基础使用
val rdd = sc.makeRDD(List(1, 2, 3, 3, 4, 4, 5, 6, 6))
val results = rdd.distinct().collect()
println(results.mkString(","))
b 去重原理分析(分布式计算的去重计算)
spark的distinct去重的核心是这个计算:
map(x => (x, null)).reduceByKey((x, _) => x, numPartitions).map(_._1)
也就是将数据源中所有数据都构成 (x, null) 的元组,然后通过reduceByKey 将所有相同值的元组聚合,最后返回该元组的第一个元素——value。
我们需要将spark和scala的distinct在实现方式上进行区分,scala是使用了HashSet进行去重操作。
总结为:
由于spark的数据是分布式的,因此采用reduceByKey将数据进行shuffle是更适合分布式场景的去重操作。
(10)coalesce算子——改变分区数量(默认用于缩减分区)

当分区数量太多而分区内数据过少,我们就可以通过coalesce算子来缩减分区,形成一个具有更少分区的RDD。
a 基础使用——缩减分区
val rdd = sc.makeRDD(List(1, 2, 3, 4, 5, 6), 6)
val coalesceRDD = rdd.coalesce(2)
coalesceRDD.saveAsTextFile("output")
b 数据倾斜 & shuffle
coalesce第二个参数是shuffle的开关,若为false(默认)则不会自动shuffle,因此数据分区缩减并不是平均分配的,这就可能造成数据倾斜。

因此我们可以将第二个参数设置为true,以在数据缩减分区时进行shuffle操作,平均分区内的数据数量。但需要注意,shuffle过程会打乱数据的顺序,比如一个RDD(List(1, 2, 3, 4, 5, 6),6)在缩减为2分区时,并不会严格按照(1,2,3)一个分区,(4,5,6)一个分区,而是会随机分配。
c 使用coalesce实现增加分区
coalesce除了缩减分区外,也是可以增加分区的。但是如果不指定shuffle操作,该算子仅会创建一个具有多余分区的RDD,但是不会进行数据落地:

因此在使用coalesce增加分区时,需要显示指定进行shuffle操作,即将第二个参数指定为true。但是每次都 指定第二个参数无疑较为繁琐,因此spark为我们提供了增加分区的简化版本,即我们下面提到的repartition算子。
(11)repartition算子——增加分区
不难发现,repartition算子底层就是调用了coalesce算子,并且指定shuffle参数 = true:

关于repartition算子的基础使用如下:
val rdd = sc.makeRDD(List(1, 2, 3, 4, 5, 6), 2)
val repartitionRDD = rdd.repartition(3)
repartitionRDD.saveAsTextFile("output")
因此,关于分区数量改变的操作总结为:
- 缩减分区:使用coalesce
- 增加分区:使用repartition
(12)sortBy算子——根据规则函数对RDD数据进行排序

该方法类似于将scala的sortBy和sortWith整合到一起,提供了一个更灵活的排序的方法。第一个参数用于指定排序的规则,第二个参数则用于指定是否降序(默认为true,即升序)。
一般来说,若我们在第一个参数中仅说明根据某个参数的数据类型进行排序,则可以结合第二个参数一起使用;若我们在第一个参数中需要对数据进行多字段分别排序,则此时需要传递一个完整规则,第二个参数可能就无需指定。这正是spark的sortBy的灵活之处!
val rdd = sc.makeRDD(List(1, 5, 2, 4, 6, 3), 2)
val sortRDD = rdd.sortBy(num => num, false) //根据num的数据类型进行降序排序
println(sortRDD.collect().mkStri(","))
(13)针对两个RDDs数据的操作——交 & 并 & 差 & 拉链
若我们想要将两个RDD的数据合并到一起,我们就需要采用类似集合的操作。注意,reduce操作仅是将单个RDD的数据进行聚合,并不是针对两个RDD的数据;同时,spark中没有fold方法,因此无法通过类似scala的fold将两个RDD进行折叠。
我们在下面要对这两个RDD进行双value相关操作:
val rdd1 = sc.makeRDD(List(1, 2, 3, 4))
val rdd2 = sc.makeRDD(List(3, 4, 5, 6))
a 并集——union

注意,关于spark中“集合”操作,针对的是seq这类数据,并不是单独针对set数据。因此,求集合后得到的RDD数据中,是可能会有重复数据的,此处提及后下文不再赘述。
//并集
val unions = rdd1.union(rdd2).collect()
println(unions.mkString(",")) // 1,2,3,4,3,4,5,6
b 交集——intersection

//交集
val intersect = rdd1.intersection(rdd2).collect()
println(intersect.mkString(",")) // 3,4
c 差集——subtract

注意,求差集的操作并不是双向的,而是 调用者 – 传入参数 的方向,因此差集元素是仅调用者有的元素。
//差集
val subtracts = rdd1.subtract(rdd2).collect()
println(subtracts.mkString(",")) // 1,2
d 拉链——zip

//拉链
val zips = rdd1.zip(rdd2).collect()
print(zips.mkString(","))
e 关于双value操作的注意事项
关于两个RDD的数据类型要求:
交、并、差都要求两个RDD数据集的数据类型相同。
但是zip操作只是将两个RDD数据集进行合并形成tuple,因此不要求两个RDD数据集数据类型相同。
关于两个RDD进行zip操作时的数据量要求:
zip操作要求被拉链的两个RDD需要保持分区数一致 & 分区内数据量一致,这与scala的zip操作是不同的。
本质上还是因为spark在分布式进行操作,scala在单机中进行操作。
(14)针对KV类型数据的操作
a partitionBy算子——按照分区规则重新分区

① partitionBy的相关说明
在使用该方法时,需要确保RDD的数据类型为KV类型,才会进行二次编译查找partitionBy方法。
其次,该方法是通过指定分区器,按照分区器规则对KV数据进行重新分配。常见分区器有以下两种:
用于按照key的Hash值进行重分区:

用于按照排序规则对数据进行分区(sortBy使用):

② partitonBy的基础使用
val rdd = sc.makeRDD(List(("a", 1), ("b", 2), ("c", 3), ("a", 2)))
rdd.partitionBy(new HashPartitioner(2)).saveAsTextFile("output")
③ partitionBy的注意事项
- 若重分区的分区器与当前RDD的分区器类型相同:

partitionBy在底层进行了分区器的判断,若使用相同分区器,说明两个分区规则相同,则返回当前RDD自身(毕竟没必要重新分区);只有当两个分区器不同的时候,才会产生新的RDD,用于存放新的分区数据。
- 若想要使用自定义分区器规则进行分区:
此时我们可以仿照源码中的分区器,写一个自定义的分区器即可。
b reduceByKey算子——按照Key进行聚合

该方法可以使我们能够在key相同的情况下,根据我们自定义的value聚合规则对value进行聚合操作。
① 基础使用
val rdd = sc.makeRDD(List(("a", 1), ("b", 2), ("c", 3), ("a", 2)))
val results = rdd.reduceByKey(_ + _).collect()
println(results.mkString(","))
② 注意事项
- 聚合操作仅会对key出现至少两次的元组进行,对于仅出现一次key的元组,不会执行相关操作。
- 同时,reduceByKey是基于scala的reduce方法,从左向右进行两两聚合 & 迭代。
c groupByKey算子——按照Key进行分组

val rdd = sc.makeRDD(List(("a", 1), ("b", 2), ("c", 3), ("a", 2)))
val results = rdd.groupByKey().collect()
println(results.mkString(","))
与groupBy不同的一点在于,该方法可以对元组数据自动进行按照Key的分组,但同时也只能够按照Key进行分组,无法像groupBy一样实现自定义分组(不过也好理解,kv的映射肯定Key是进行分组的依据,不然放到Key的位置干啥)。
d 区分:reduceByKey & groupByKey的区别和使用场景
首先需要知道,spark的shuffle操作,都是会进行落盘的。因为shuffle后的操作是需要分区内全部数据的,否则就会造成数据的缺漏(比如a, b, c分别在三个分区,在shuffle后若不等三个数据全部就位,直接对a, b进行计算,则c数据将会被遗漏);同时,等待数据全部到位的操作若在内存中执行,在数据量大的情况下会造成内存溢出。因此,对于shuffle操作,采用写入磁盘,后续操作再从磁盘中读取数据进行处理。

reduceByKey和groupByKey,都是根据元组数据的Key进行操作,且都是shuffle操作,同样会进行落盘。我们能够先groupByKey,再通过map进行聚合,何必使用reduceByKey的操作呢?
核心原因在于,reduceByKey会在shuffle前,先对分区内相同Key的数据进行reduce,然后再将shuffle后的数据按照Key进行reduce。也就是说,reduceByKey可以减少shuffle的数据数量,减少了IO的交互,这显然提升了性能。

但同时,我们需要明确groupByKey和reduceByKey的设计理念:groupByKey用于将数据进行分组;reduceByKey用于将数据进行按Key聚合。二者没有优劣之分,只有使用场景上的差异!
e aggregateByKey算子——对分区内和分区间分别指定计算规则

① 形参说明
该方法有三个参数(严格上来说是两个参数,因为第二个参数包含两个函数参数),我们先对三个参数进行说明:
- zeroValue:初始值。与reduceByKey不同,该方法每次仅获取一个分区数据,依据分区数据与初始值的计算结果作为下一次计算的初始值,不断迭代计算(参考scala的fold)。
- seqOp:分区内计算函数。该参数用于指定分区内以什么规则进行计算。分区内计算实际上就是不断将计算结果作为下一个初始值,与下一个分区数据进行计算。注意,这个计算函数有两个参数,第一个参数就是进行迭代的初始值,第二个参数则是当前的value。
- combOp:分区间计算函数。该参数用于指定分区聚合时的计算规则,两个参数是不同分区的相同key数据。
② 实操——分区内按Key求最大值,分区间按Key进行求和
我们在5.2.7-(5)glom算子中曾使用glom + map 进行操作,这里我们使用aggregateByKey一步到位,直接指定分区内 & 分区间的计算规则:
val rdd = sc.makeRDD(List(("a", 1), ("a", 2), ("a", 2), ("a", 3)), 2)
//取出每个分区内相同key的最大值,然后进行分区间求和
// 初始值为了碰上第一个key,进行分区内计算,因为此时无法获取分区内第二个元素
val results = rdd.aggregateByKey(0)(
(x, y) => math.max(x, y), //分区内计算规则,返回值作为下一个初始值参与分区内计算
(x, y) => x + y //分区间计算规则
).collect()
println(results.mkString(","))
该案例的分区内计算规则中,x为迭代的初始值(每一次比较后的max),y为当前进行比较的元素。
分区间计算规则中,x 和 y为不同分区相同key的最大值。
③ 对aggregateByKey的深入理解 & 案例实现
我们在形参说明的时候提到,分区内对value的计算是通过对初始值的迭代进行计算的,因此返回结果需要与初始值的类型保持一致。同时,分区间的计算,也是对分区内计算后得到的value进行聚合,所以分区间计算的返回值也应该与初始值类型保持一致(value的类型)。最后形成的将是一个 key -> 初始值类型 的元组数据集。
现在我们根据这个概念,就可以对初始值类型进行灵活的设置,以满足我们的需求。比如我们现在要求每个Key的value均值:
首先明确需求:均值 = value和 / key的出现次数,因此我们要获取的值有两个,也就是∑value & key出现次数。
基于需求,我们可以设置一个tuple类型的初始值,这个初始值的含义为 (value的和, key出现的次数)。我们在分区内每一次获取相同的key,都可以对这个初始值进行value求和 & 次数增加的操作;分区内计算完成后,我们就得到了分区内每个key对应的 (∑value, 出现次数),然后在分区间对相同key的value再进行聚合即可:
val results = rdd
//求出每个key的value和 & key出现次数
.aggregateByKey((0, 0))( //初始值含义:第一个参数为value和,第二个参数为出现次数
(tuple, value) => (tuple._1 + value, tuple._2 + 1), //tuple为迭代初始值,进行当前key两个参数的计算;value为当前key对应的value
(tuple1, tuple2) => (tuple1._1 + tuple2._1, tuple1._2 + tuple2._2) //tuple1 和tuple2为当前key在不同分区的value
)
//对value进行平均值的map(key相同)
.mapValues(value => value._1 / value._2)
.collect()
println(results.mkString(“,”))
f foldByKey算子——分区间和分区内计算规则相同的简化

由于是对aggregateByKey的简化,所以该方法也需要一个初始值,来对第一个Key进行计算。
val rdd = sc.makeRDD(List(("a", 1), ("b", 2), ("b", 2), ("a", 3)), 2)
//分区内按key对value求和,分区间也按key对value求和
val results = rdd.foldByKey(0)(_ + _).collect()
println(results.mkString(","))
g combineByKey算子——直接将第一个数据转换作为初始值进行计算
aggregateByKey 和 foldByKey 都能够基于初始值对分区内数据进行计算,但是这个初始值是我们指定的,与原始数据并没有关系,只是在运算时与数据进行运算。
combineByKey算子可以直接将分区内相同key的第一个数据转换为我们自定义的初始值格式,然后通过该初始值与后续相同key的value进行我们自定义的分区内 & 分区间计算:

combineByKey的参数说明如下:
- createCombiner:将出现的第一个key的数据转换为初始值,这个初始值我们可以自定义格式。
- mergeValue:指定在分区内将初始值与当前key的value进行聚合计算的函数。
- mergeCombiners:指定在分区间将分区内计算结果基于相同key进行聚合计算的函数。
基于此我们就可以得出另一个求平均值的实现方式:
val rdd = sc.makeRDD(List(("b", 4), ("a", 5), ("b", 2), ("a", 3)), 2)
//求出每个key的value均值
val results = rdd
.combineByKey(
value => (value, 1), //将相同key的第一个元素转换为自定义格式的初始值:(∑value, key出现的次数)
(tuple: (Int, Int), value) => (tuple._1 + value, tuple._2 + 1), //分区内对相同key的value计算规则
(t1: (Int, Int), t2: (Int, Int)) => (t1._1 + t2._1, t1._2 + t2._2) //分区间对相同key的value聚合计算规则
)
.mapValues(t => t._1 / t._2)
.collect()
println(results.mkString(","))
h 区分:reduceByKey & aggregateByKey & foldByKey & combineByKey的区别以及使用场景
- reduceByKey:直接将相同key的value进行两两计算,并且分区内和分区间计算规则相同。
- foldByKey:自定义初始值,在分区内计算时将相同key和初始值进行两两计算,然后不断迭代,分区内和分区间计算规则相同。
- aggregateByKey:自定义初始值,每个key使用的初始值相同,在分区内计算时将相同key和初始值进行两两计算,然后不断迭代,分区内和分区间计算规则不同。
- combineByKey:将相同key的第一个数据直接转换为我们需要的初始值,初始值动态变化(与key相关),然后与相同key数据进行计算,分区内和分区间计算规则不同。
i join & leftOuterJoin & rightOuterJoin算子——连接操作
该操作与数据库表中的连接操作基本类似。与zip不同,该操作是对key进行连接,而zip操作是基于RDD中元素的索引进行连接。
① join算子

只对两个RDD均有的Key进行连接。
val rdd1 = sc.makeRDD(List(("a", 1), ("b", 3), ("c", 4)))
val rdd2 = sc.makeRDD(List(("b", 2), ("a", 1), ("d", 4)))
val results = rdd1.join(rdd2).collect()
println(results.mkString(",")) //(a,(1,1)),(b,(3,2))
② leftOuterJoin & rightOuterJoin算子
leftOuterJoin:将调用者作为主表进行连接,主表数据全参与。

val rdd1 = sc.makeRDD(List(("a", 1), ("b", 3), ("c", 4)))
val rdd2 = sc.makeRDD(List(("b", 2), ("a", 1), ("d", 4)))
val results = rdd1.leftOuterJoin(rdd2).collect()
println(results.mkString(",")) //(a,(1,Some(1))),(b,(3,Some(2))),(c,(4,None))
rightOuterJoin:将参数作为主表进行连接。

val rdd1 = sc.makeRDD(List(("a", 1), ("b", 3), ("c", 4)))
val rdd2 = sc.makeRDD(List(("b", 2), ("a", 1), ("d", 4)))
val results = rdd1.rightOuterJoin(rdd2).collect()
println(results.mkString(",")) //(d,(None,4)),(a,(Some(1),1)),(b,(Some(3),2))
j cogroup算子——不同RDD内相同key的分别聚合

这个算子能将两个RDD的相同key分别聚合成Iterator,然后将两个迭代器进行连接,获得一个相同key,不同RDD的value集合。下图中,我们用rdd1.cogroup(rdd2),那么就会将rdd1和rdd2中所有key分别group到一起(比如对于key = a,rdd1中形成(a, CompactBuffer(1, 2, 3)),rdd2中形成(a, CompactBuffer(3))),然后将两个相同key进行连接,最终形成(key, rdd1中的CompactBuffer, rdd2中的CompactBuffer) 这样的结构:
val rdd1 = sc.makeRDD(List(("a", 1), ("b", 2), ("a", 2), ("a", 3)))
val rdd2 = sc.makeRDD(List(("a", 3), ("b", 4), ("b", 2), ("c", 2)))
val results = rdd1.cogroup(rdd2).collect()
results.foreach(println)
//(a,(CompactBuffer(1, 2, 3),CompactBuffer(3)))
//(b,(CompactBuffer(2),CompactBuffer(4, 2)))
//(c,(CompactBuffer(),CompactBuffer(2)))
(15)转换算子案例实操

a 数据源部分数据如下:

b 实现思路:
- 获取原始数据,将其按空格拆分,形成 (时间戳, 省份, 城市, 广告, 用户) 类型的数据。
- 将拆分后数据选出 (省份,广告) ,并且对每一条数据都map成 ((省份, 广告), 1) 的形式。这个1表示当前广告的出现次数,我们对每一条数据都采用1作为默认值,便于我们后续进行求和统计。
- 按照第二步中的数据形式,按key对value进行求和,统计这个广告在这个省内出现的次数,将数据map成 ((省份, 广告), 出现次数) 。
- 将统计求和后的数据进行转换,map成 (省份, (广告, 出现次数)) 的形式,让省份成为唯一key,便于后续基于省份划分广告。
- 对第三步中的数据按省份进行分组,并且进行map操作,对value进行按照出现次数的降序排列,取出value的前三,构成 (省份 -> 前三广告集合) 的数据类型,该数据即为结果。
c 代码实现:
val sparkConf = new SparkConf().setMaster("local[*]").setAppName("agent.log")
val sc = new SparkContext(sparkConf)
//读取数据
val rdd = sc.textFile("datas/agent.log")
//将数据进行切分
val rdd1 = rdd.map(_.split(" "))
//按省份 & 广告进行数据转换,构成(省份,广告,出现次数)的数据结构
val rdd2 = rdd1.map(line => ((line(1), line(4)), 1)) //默认出现次数为1
//对这个(省份,广告)的出现次数进行统计
val rdd3 = rdd2.reduceByKey(_ + _)
//将数据再转换为(省份,(广告,出现次数)),使得省份成为唯一的key
val rdd4 = rdd3.map(kv => (kv._1._1, (kv._1._2, kv._2)))
//按照省份进行再次分组,对每组的值都按广告出现次数(即_._2)进行降序排列,选出每组的前三
val rdd5 = rdd4.groupByKey().map(kv => kv._1 -> kv._2.toList.sortWith(_._2 > _._2).take(3))
val results = rdd5.collect()
results.foreach(println)
//(4,List((12,25), (2,22), (16,22)))
//(8,List((2,27), (20,23), (11,22)))
//(6,List((16,23), (24,21), (22,20)))
//(0,List((2,29), (24,25), (26,24)))
//(2,List((6,24), (21,23), (29,20)))
//(7,List((16,26), (26,25), (1,23)))
//(5,List((14,26), (21,21), (12,21)))
//(9,List((1,31), (28,21), (0,20)))
//(3,List((14,28), (28,27), (22,25)))
//(1,List((3,25), (6,23), (5,22)))
5.2.8 RDD基础编程——行动算子

所谓行动算子,就是触发 spark job 执行的方法,代码底层会调用环境对象的 runJob 方法,根据依赖关系创建DAG,创建ActiveJob,提交并执行。
与转换算子不同,行动算子返回结果不再是新RDD,而是具体数据类型。
(1)reduce算子——将RDD数据聚合并返回

该算子与先前的reduceByKey不同,直接返回了具体的数据类型(即调用者RDD中的数据类型)。
val rdd = sc.makeRDD(List(1, 2, 3, 4))
//reduce
val sum = rdd.reduce(_ + _)
println(sum) //10
(2)collect算子——将RDD数据收集到内存形成结果数组

该方法会将RDD数据按分区顺序存入Driver的内存中形成数组,同时,为了防止内存溢出,这个结果数组不应过大(这也是该方法给我们的note)。
val rdd = sc.makeRDD(List(1, 2, 3, 4))
//collect——按分区收集
val collect = rdd.collect()
println(collect.mkString(",")) //1,2,3,4
(3)count算子——计算RDD中的数据个数

val rdd = sc.makeRDD(List(1, 2, 3, 4))
//count
val nums = rdd.count()
println(nums) //4
(4)first算子——取出RDD中第一个数据

val rdd = sc.makeRDD(List(1, 2, 3, 4))
//first
val first = rdd.first()
println(first) //1
(5)take算子——取出RDD中前n个数据形成Array

val rdd = sc.makeRDD(List(1, 2, 3, 4))
//take
val takes = rdd.take(3)
println(takes.mkString(",")) //1,2,3
(6)takeOrdered算子——返回RDD中按照排序规则排序后的数据的前n个

//takeOrdered——返回按照排序后的前n个数据
val rdd1 = sc.makeRDD(List(1, 8, 6, 4))
val ordered = rdd1.takeOrdered(2)
println(ordered.mkString(",")) //1,4
(7)aggregate算子——分别与初始值进行分区内 & 分区间运算

与aggregateByKey算子不同,该算子的初始值会同时参与分区内 & 分区间的计算。
val rdd = sc.makeRDD(List(1, 2, 3, 4), 2)
val results = rdd.aggregate(4)(_ + _, _ * _)
//分区内与初始值求和,分区间与初始值进行相乘(分区内 & 分区间运算全参与)
println(results) //308
上图代码中,初始值4参与了分区内的 _ + _ 运算 以及分区间的 _ * _ 运算,因此实际运算应为 4[初始值] * (4[初始值] + 1 + 2) * (4[初始值] + 3 + 4) = 308。
(8)fold算子——aggregate算子的简化,分区内 & 分区间计算规则相同

同样的,该算子的初始值也会参与分区内 & 分区间的计算。
val rdd = sc.makeRDD(List(1, 2, 3, 4), 2)
val results1 = rdd.fold(4)(_ + _)
println(results1) //22
计算过程为: 4[初始值] + (4[初始值] + 1 + 2) + (4[初始值] + 3 + 4) = 22。
(9)countByValue & countByKey算子——统计value出现的次数& 统计Key出现的次数
a countByValue算子

该算子统计的是RDD中数据出现的次数,也就是每一个RDD数据都被视作一个value,不管这个数据是什么类型的。
注意,该算子并不是针对KV的V进行统计,若一个具有KV类型的RDD调用该方法,则该方法会将整个KV视作一个元素,而不是按tuple的KV统计!
val rdd1 = sc.makeRDD(List(1, 2, 3, 4), 2)
val valueCount1 = rdd1.countByValue()
println(valueCount1) //Map(4 -> 1, 2 -> 1, 1 -> 1, 3 -> 1)
//KV数据将被视作一个整体
val rdd2 = sc.makeRDD(List(("a", 1), ("b", 2), ("c", 1)), 2)
val valueCount2 = rdd2.countByValue()
println(valueCount2) //Map((b,2) -> 1, (a,1) -> 1, (c,1) -> 1)
b countByKey算子

该算子才是针对KV数据进行统计的算子,统计结果为Key出现的次数,与value无关。
val rdd2 = sc.makeRDD(List(("a", 1), ("b", 2), ("c", 1)), 2)
val keyCount = rdd2.countByKey()
println(keyCount) //Map(b -> 1, a -> 1, c -> 1)
(10)save相关算子——保存结果到文件中
a saveAsTextFile——保存为text文件:

b saveAsObJectFile——将数据序列化为对象保存到文件中:

c saveAsSequenceFile——将KV类型转换为Sequence保存到文件中:

d 使用展示
val rdd1 = sc.makeRDD(List(1, 2, 3, 4))
rdd1.saveAsTextFile("output1") //保存为text文件
rdd1.saveAsObjectFile("output2") //序列化为对象保存到文件
val rdd2 = sc.makeRDD(List(("a", 1), ("b", 2), ("c", 3)))
rdd2.saveAsSequenceFile("output3") //保存成SequenceFile文件(仅KV类型可用)
(11)foreach算子——分布式遍历

RDD的foreach算子是分布式遍历,也就是将RDD数据分发到不同的Executor后进行遍历,因此数据遍历顺序会根据不同分区的空闲状态而不同,并不固定:
val rdd = sc.makeRDD(List(1, 2, 3, 4))
//直接调用foreach算子
rdd.foreach(println(_)) //3, 1, 4, 2(无序)
而我们之前使用的对集合的foreach是Scala提供的方法,直接对保存后的集合进行遍历,注意,这个遍历并不是分布式的:
val rdd = sc.makeRDD(List(1, 2, 3, 4))
//保存为Array后再进行foreach遍历(调用scala的集合方法)
rdd.collect().foreach(println) //1, 2, 3, 4
我们可以在这里再次明确算子的概念:算子的概念是为了区分RDD与scala集合操作,scala集合的所有操作都是在同一片内存中执行的;而RDD算子属于分布式执行,会分散到不同的Executor中执行。
我们可以理解为:
- Scala的集合操作等算子以外的操作是在Driver内存中执行的。
- Spark RDD的算子是Driver分配后在不同Executor内存中分布执行的。
5.2.9 序列化
(1)引入——从foreach算子遍历对象到序列化
现在,我们想用RDD算子foreach去遍历一个Student对象:
class Student(var name: String, var age: Int)
val rdd = sc.makeRDD(List(1, 2, 3, 4))
val student = new Student("张三", 19)
rdd.foreach(num => println(s"num : ${num}, name : ${student.name}, age : ${student.age}"))
结果我们会发现,这个遍历会报错:

大致意思就是说:这个foreach任务没有序列化,原因是Student这个对象没有进行序列化。我们将Student混入Serializable特质(或者直接将其用case定义为样例类),就可以解决这个问题:
//将该类混入序列化特质,使其可以被序列化
class Student(var name: String, var age: Int) extends Serializable
val rdd = sc.makeRDD(List(1, 2, 3, 4))
val student = new Student("张三", 19)
rdd.foreach(num => println(s"num : ${num}, name : ${student.name}, age : ${student.age}"))

那么,为什么需要序列化对象,才能使用RDD算子进行计算呢?
首先,在网络传输上,传输内容需要可以序列化为对象才能去传输,这是所有网络传输的基础。然后在spark计算框架中,RDD是由Driver基于分区、计算任务等将任务分发给Executor去执行的,这个过程是需要进行网络传输的。由于RDD的分配是一个闭包,分发给Executor的计算任务需要知道他上一步的数据,所以闭包内的数据一定要可以被序列化,才能进行传输。
(2)RDD的闭包检测(Driver端执行)
在Driver端准备分发任务时,会调用clean方法来清理闭包:

该步骤的目的在于 清理不必要的引用(闭包不会用上的数据) & 检查序列化的可行性。对于序列化可行性的检查,主要逻辑在于如下代码:

若检测到的对象无法进行序列化,那么该闭包就无法被正确传输,因此无法将任务分发给Executor去执行。
(3)Kryo序列化(大数据场景下序列化的改善)
参考地址: https://github.com/EsotericSoftware/kryo
Java的序列化能够序列化任何的类。但是比较重(字节多),序列化后,对象的提交也 比较大。
Spark出于性能的考虑,Spark2.0开始支持另外一种Kryo序列化机制。Kryo速度 是Serializable的10倍。当RDD在Shuffle数据的时候,简单数据类型、数组和字符串类型 已经在Spark内部使用Kryo来序列化。
object serializable_Kryo {
def main(args: Array[String]): Unit = {
val conf: SparkConf = new SparkConf()
.setAppName("SerDemo")
.setMaster("local[*]")
// 替换默认的序列化机制
.set("spark.serializer",
"org.apache.spark.serializer.KryoSerializer")
// 注册需要使用 kryo 序列化的自定义类
.registerKryoClasses(Array(classOf[Searcher]))
val sc = new SparkContext(conf)
val rdd: RDD[String] = sc.makeRDD(Array("hello world", "hello atguigu",
"atguigu", "hahah"), 2)
val searcher = new Searcher("hello")
val result: RDD[String] = searcher.getMatchedRDD1(rdd)
result.collect.foreach(println)
}
}
case class Searcher(val query: String) {
def isMatch(s: String) = {
s.contains(query)
}
def getMatchedRDD1(rdd: RDD[String]) = {
rdd.filter(isMatch)
}
def getMatchedRDD2(rdd: RDD[String]) = {
val q = query
rdd.filter(_.contains(q))
}
}
5.2.10 依赖关系

简单理解:
依赖关系指的是两个RDD之间的关系
血缘关系指的是多个连续的RDD之间的依赖关系
(1)血缘关系

我们刚才提到了,血缘关系就是一系列RDD之间的依赖关系。由于RDD是不保存数据的,所以RDD会将血缘关系保留下来,以提供容错性。
我们针对下图进行保留血缘关系的分析:

- RDD4的血缘关系构成:RDD4会从数据源中通过textFile读取datas/word.txt,这将会作为一个血缘保留在RDD4中,告诉RDD4如何获取本次计算的数据。这也是为什么RDD不保存数据却能够获取数据的原因,它会保留数据的获取途径。
- RDD3的血缘关系构成:RDD3会从RDD4中读取数据,并且进行flatMap操作,因此会将RDD4的读取文件方式和文件路径保留在其血缘关系中,然后再将flatMap操作添加到这个血缘关系中。
- RDD2的血缘关系构成:RDD2会从RDD3中获取操作数据,进行map操作,也就是依赖RDD3的数据,因此会将RDD3中的血缘关系(即RDD4的读取文件 & RDD3的flatMap)保留在RDD2中,然后再将map操作添加到RDD2的血缘关系中。
- RDD1的血缘关系构成:最后,RDD1会从RDD2中获取数据,进行reduceByKey操作,因此会将RDD2的血缘关系(textFile + flatMap + map)保留在其血缘关系中,再添加reduceByKey操作到RDD1的血缘关系中。
至此,RDD4 -> RDD1的血缘关系记录完毕,且每一个RDD都包含自身的血缘关系。当出现某一步的执行出错,就可以依据这个血缘关系,重新执行操作。且这个操作无需重新获取先前的RDD,因为血缘关系保留在自身RDD上。
(2)宽窄依赖
宽窄依赖是基于RDD分区数据在新旧RDD间的使用情况来划分的,针对的是 父RDD一个分区的数据 在 子RDD分区内的使用情况 。若父子之间分区是一对一的使用,那么则称为窄依赖;若父子之间分区数据是一对多使用,那么则称为宽依赖。
宽窄依赖是RDD是否进行shuffle操作的重要标识,也是RDD进行阶段划分的重要依据,是RDD对MR每一步都落盘操作的优化基石。
a 窄依赖(OneToOne依赖)
窄依赖表示每一个父(上游)RDD的Partition最多被子(下游)RDD的一个Partition使用:

b 宽依赖(Shuffle依赖)
宽依赖表示同一个父(上游)RDD的Partition被多个子(下游)RDD的Partition依赖,会 引起Shuffle:

(3)RDD阶段划分
a 概念解析
① 窄依赖时的阶段划分
若新RDD的分区只需要原RDD分区的数据就可以执行计算(分区之间一对一,类似管道传输),那么阶段就无需划分,两个RDD的任务可以在同一阶段执行,且两个分区之间的任务并不会互相影响。
比如下图,新旧RDD之间的分区数据使用是一对一的,上方分区只使用上方分区数据,下方分区只使用下方分区数据,因此二者的执行互不干扰,可以将上方分区划分为一个任务,下方分区划分为一个任务,二者在同一阶段一起执行。

② 宽依赖时的阶段划分
但是若新RDD的数据需要旧RDD中多个分区的数据,那么新旧RDD的计算任务就无法在同一个阶段执行完成了。因为新RDD必须等旧RDD所有分区数据就绪,否则数据将会出现遗漏。那么此时,新旧RDD的计算任务就会被划分到不同的阶段。
如下图, 新RDD上方分区 需要 旧RDD两个分区的数据 ,下方分区同理。此时,若将 新旧RDD上方分区 统一划分为一个任务就不合理了,这样就无法接收下方分区的数据。因此, 旧RDD的两个分区 和 新RDD的两个分区计算 需要划分不同的任务,同时, 旧RDD的计算 和 新RDD的计算 需要划分到两个不同阶段。

b 源码分析
① 阶段划分源码的调用情况
- 在RDD执行行动算子时,会调用runJob方法,该方法会调用submitJob,提交要执行的任务:

- 在submitJob中,会提交JobSubmitted事件:

- 对于JobSubmitted事件,由handleJobSubmitted方法去处理这个事件,在handleJobSubmitted中,就会调用阶段划分的方法createResultStage,这是阶段划分的逻辑代码:


② 阶段划分的具体执行
阶段划分简单来说就是按照shuffle依赖来进行划分,接着我们来看看具体是如何进行阶段划分,以将所有shuffle划分到不同阶段的。
- 在我们最开始调用阶段划分时,此时处于啥也没有的阶段,因此createResultStage会直接创建一个stage = new ResultStage(其实这一步是在得到父阶段、任务id等属性后完成的,但是我们可以将其视作先创建一个空白ResultStage,然后将获取的属性添加进去):

- 同时,由于当前阶段需要获取上一阶段的数据,因此会通过getOrCreateParentStages去查找 血缘关系 中的父阶段:

- 在通过getOrCreateParentStages获取父阶段中,会通过查找是否有shuffle依赖来判断父阶段(我们之前提到过,阶段是按照宽依赖划分的,即否进行shuffle操作):

- 在getShuffleDependencies方法中,会通过模式匹配,判断当前获取的依赖是否是shuffle依赖,若是,则将其添加到父阶段parents中,最后将获取到的所有父阶段(shuffle依赖)返回给getOrCreateParentStages方法,通过其去进行整合:

- 在得到了所有父阶段parents(shuffle依赖)后,getOrCreateParentStages就会对所有的shuffle依赖进行映射,对每一个shuffle依赖都通过getOrCreateShuffleMapStage创建一个新阶段,将结果通过toList整合起来,作为parents返回给getOrCreateParentStages,使其作为属性能够添加到当前阶段stage中:


- 最后,父阶段划分完毕,同时将当前阶段(实际上是在new ResultStage完成的)赋予 jobId、父阶段、分区等内容,将这个stage返回给handleJobSubmitted,作为此次需要执行的阶段:

③ 阶段划分的总结
简单来说,阶段划分就是:创建新阶段,即:赋予新阶段任务 + 找父阶段(shuffle依赖),最后将这个阶段返回。

- 首先是阶段划分,只要有一个shuffle依赖,那么就会自动增加一个阶段。
- 其次关于阶段的数量,由于shuffle依赖是属于父阶段,在现有阶段已经通过createResultStage创建了一个新的阶段,因此总阶段数量 = 父阶段数量(shuffle依赖数量) + 1(当前创建的新阶段ResultStage)。
- 最后,关于我们创建的ResultStage,这是我们当前阶段需要执行的任务,也就是所有父阶段执行完后,才需要执行的任务,因此会在最后执行。
(4)RDD任务划分
a 任务划分
RDD任务切分中间分为:Application(应用)、Job(工作)、Stage(阶段)和Task(任务)。
- Application(应用):初始化一个SparkContext即生成一个Application;
- Job(工作):一个Action算子就会生成一个Job;
- Stage(阶段):Stage 等于宽依赖(ShuffleDependency)的个数加1;
- Task(任务):一个Stage阶段中,最后一个RDD的分区个数就是Task的个数。
注意:Application -> Job -> Stage -> Task 每一层都是 1 对 n 的关系:
- 一个Application中可能有多个行动算子,构成多个Job
- 一个Job中可能有多次shuffle算子,构成多个Stage
- 一个stage中最后的RDD可能有多个分区,构成多个Task

b 任务划分的源码解读
- 在handleJobSubmitted方法中,除了会对阶段进行划分,同时也会提交阶段,进行阶段内任务的划分:

- 任务划分是依据阶段的类型(实际上,由于shuffle后的一个阶段内,都是OneToOne依赖,因此二者任务划分基于同样的过程),基于partitionsToCompute进行任务划分,也就是根据分区进行任务计算:

- 在partitionsToCompute中,会基于当前stage调用findMissingPartitions,来进行任务计算:

- 下面是ShuffleStage的findMissingPartitions计算,实际上ResultStage调用的是同样的过程,会根据当前Stage中最后RDD的分区数来划分,有几个分区也就是有几个任务:

总结来说,就是阶段内RDD有几个分区,那么就划分为几个任务。此处不用考虑分区增减的情况,因为若分区增减涉及shuffle,那就属于下一个阶段;若分区不涉及shuffle,那么只是单纯增加空白分区,RDD间仍是OneToOne依赖,进行计算的分区不变,任务划分也不用改变。
5.2.11 持久化
(1)持久化的引入
假设我们正在对list中的数据进行wordCount & 针对Key进行分组:

那么上述代码的流程可以简化为下图的流程:

我们会发现,rdd3(reduceByKey) & rdd4(groupByKey)都使用了rdd2提供的结果数据。但是由于RDD仅保留血缘关系,不保留数据,因此两个rdd的操作之间数据无法共享,也就是所有操作均需要重新再来一遍。
那么我们为了数据共享,提升相同计算rdd的效率,我们可以在map后进行持久化操作,将数据持久化到内存 / 磁盘中,让reduceByKey & groupByKey两个算子可以直接从共享数据域中读取数据,这样就可以直接执行这两个操作,无需重新计算:

需要注意,持久化的位置是需要基于需求考虑的,若要数据安全性则落盘会好一点,若要计算的效率则内存会好一点。
(2)在代码中进行持久化

我们只需要将想要共享的RDD通过.cache() 或 .persist() 就可以将RDD数据持久化到内存 / 磁盘。
实际上,.cache()底层是调用了persist方法,将存储级别指定为内存;而若想要自行指定存储级别,可以使用persist方法并为其指定参数:

存储级别如下:


(3)持久化的作用总结
RDD通过Cache或者Persist方法将前面的计算结果缓存,默认情况下会把数据以缓存 在JVM的堆内存中。但是并不是这两个方法被调用时立即缓存,而是触发后面的action算 子时,该RDD将会被缓存在计算节点的内存中,并供后面重用。
缓存有可能丢失,或者存储于内存的数据由于内存不足而被删除,RDD的缓存容错机 制保证了即使缓存丢失也能保证计算的正确执行。通过基于RDD的一系列转换,丢失的数 据会被重算,由于RDD的各个Partition是相对独立的,因此只需要计算丢失的部分即可, 并不需要重算全部Partition。

同时,持久化不一定只能用于数据重用。对于一些 耗时较久的任务 或 数据比较重要的场合也可以使用。Spark 会自动对一些Shuffle操作的中间数据做持久化操作(比如:reduceByKey)。这样 做的目的是为了当一个节点Shuffle失败了避免重新计算整个输入。但是,在实际使用的时 候,如果想重用数据,仍然建议调用persist或cache。
(4)CheckPoint检查点
所谓的检查点其实就是通过将RDD中间结果写入磁盘。
由于血缘依赖过长会造成容错成本过高,这样就不如在中间阶段做检查点容错,如果检查点 之后有节点出现问题,可以从检查点开始重做血缘,减少了开销。
同时,对RDD进行checkpoint操作并不会马上被执行,必须执行Action操作才能触发,这与持久化触发要求是一致的。

检查点文件会保存到我们刚刚设置的目录(基于根目录的相对路径)下:

注意,一般而言,CheckPoint要存储与hdfs中。
(5)持久化 & CheckPoint 的区别
cache:
- 将数据临时存储在内存中进行数据重用。进程结束,则内存中数据丢失。
- 会在血缘关系中添加新的依赖(告诉血缘有一个缓存操作),一旦出现问题就可以从头读取数据。
persist:
- 将数据临时存储在磁盘文件中进行数据重用。
- 过程涉及磁盘IO,性能降低但是同时数据更加安全。
- 若作业执行完毕,保存的数据文件会直接删除。
checkpoint:
- 将数据长久保存在磁盘文件中进行数据重用。
- 涉及磁盘IO,性能低但是数据安全。
- 一般而言,checkpoint作业会独立执行,以保证数据安全。
- 同时,为了提升效率,checkpoint会配合cache一起使用,既保证数据安全,又能提升性能。
- 执行的过程中,checkpoint会切断血缘关系,从checkpoint开始建立新的血缘关系。因此我们可以理解为,checkpoint改变了数据源,让后续rdd从checkpointRDD开始读取数据。
5.2.12 分区器
(1)Hash分区器

按照key的hash值决定数据的分区。
(2)Range分区器

将数据均匀映射到分区中,且每个分区数据都是有序的。
(3)自定义分区器的实现案例
现在有一种情况,对于下面这个RDD中的数据,我们想要让key = “nba”在一个分区,key = “wnba”在一个分区,其余则统一在一个分区:

由于”nba”和”wnba”以及其余key并不能简单通过hash值进行划分,因此我们需要自定义一个分区器,以实现我们的需求。
我们可以参考一下HashPartitioner是如何实现的:

继承Partitioner,在主构造器设置partitions参数作为分区数量,在getPartition方法中根据key对数据进行分区决定。
若我们要自定义实现分区器,则同样需要继承Partitioner:

两个参数的意义分别为:分区器创建的分区数量 & 按key分区的依据。
因此我们就可以基于此创建我们自己的分区器了(对分区数量的指定采用主构造器传参):
/**
* 自定义分区器
* 1、继承Partitioner
* 2、numPartitions:分区数量
* 3、getPartition:根据key决定数据的分区去向
*/
class MyPartitioner(var partitions: Int) extends Partitioner {
override def numPartitions: Int = partitions
override def getPartition(key: Any): Int = key match {
case "nba" => 0
case "wnba" => 1
case _ => 2
}
}
在使用时,我们只需要将这个分区器创建并且指定相关参数即可:
val rdd1 = sc.makeRDD(List(
("nba", "xxxxxxxx"),
("wnba", "xxxxxxxx"),
("cba", "xxxxxxxxxx"),
("nmd", "xxxxxxxxxx")
))
//自定义分区器的使用,通过partitionBy + MyPartitioner进行
val partRDD = rdd1.partitionBy(new MyPartitioner(3))
partRDD.saveAsTextFile("part")
这是最终分区文件的结果:


5.2.13 文件的读取 & 保存(保存具体源码参考5.2.8-save相关算子)
Spark 的数据读取及数据保存可以从两个维度来作区分:文件格式以及文件系统。
文件格式分为:text文件、csv文件、sequence文件以及Object文件;
文件系统分为:本地文件系统、HDFS、HBASE以及数据库。
(1)text文件读取 & 保存
// 读取输入文件
val inputRDD: RDD[String] = sc.textFile("input/1.txt")
// 保存数据
inputRDD.saveAsTextFile("output")
(2)sequence文件读取 & 保存
SequenceFile 文件是Hadoop 用来存储二进制形式的key-value对而设计的一种平面文件(Flat File)。在 SparkContext 中,可以调用sequenceFile[keyClass, valueClass](path)来读取文件。
// 保存数据为
SequenceFile dataRDD.saveAsSequenceFile("output")
// 读取SequenceFile文件
sc.sequenceFile[Int,Int]("output").collect().foreach(println)
(3)object文件读取 & 保存
对象文件是将对象序列化后保存的文件,采用Java的序列化机制。可以通过objectFile[T: ClassTag](path)函数接收一个路径,读取对象文件,返回对应的RDD,也可以通过调用 saveAsObjectFile()实现对对象文件的输出。因为是序列化所以要指定类型。
// 保存数据
dataRDD.saveAsObjectFile("output")
// 读取数据
sc.objectFile[Int]("output").collect().foreach(println)
5.3 数据结构——累加器 & 广播变量
为了避免shuffle操作,spark提供了累加器 & 广播变量,以在分布式计算的情况下提升性能。
累加器:分布式只写变量
广播变量:分布式只读变量
5.3.1 累加器(不用shuffle的分布式聚合)
(1)引入——没有累加器时分布式计算的问题
val rdd = sc.makeRDD(List(1, 2, 3, 4))
//分区内 & 分区间计算
//val i = rdd.reduce(_ + _)
var sum = 0
rdd.foreach(sum += _)
println(sum)
由于reduce操作涉及分区内 & 分区间的计算,是需要进行shuffle的,因此我们想要通过在Driver中设置一个变量,在算子中对该变量进行求和,以实现遍历中对变量的累加,不进行shuffle操作。
在上述job中,我们将sum作为求和变量,不断与RDD内的数据进行求和,期望最终返回RDD数据和。但实际输出却与我们想象的完全不符:

这是由于,虽然sum作为闭包参数被传递给了Executor端,但是Executor端对sum进行参数累加时,Driver端的sum值并没有被改变。同时,Executor端的sum也不会返回给Driver端(代码中没有这个要求)。因此造成的结果是:Executor端使用的sum与Driver端的sum毫无关系,在最终我们要在Driver端输出sum时,也就只能输出其一开始的定义0。我们可以结合下图来理解上述过程:

(2)累加器的原理 & 简单使用
在(1)中,我们讨论了sum变量仅能作为闭包传递给Executor,但无法返回给Driver。那么,假设这个sum作为一个累加器,能够将Executor执行后的结果返回给Driver,再由Driver端对该变量进行聚合。这个累加器是一个分布式共享的,Driver的累加器在Executor中可用,全局都使用同一个累加器。通过设置累加器,就能实现我们想要的结果:

累加器的作用即:
用来把Executor端变量信息聚合到Driver端。在Driver程序中定义的变量,在 Executor端的每个Task都会得到这个变量的一份新的副本,每个task更新这些副本的值后, 传回Driver端进行merge。
具体而言,Driver端可以设置一个累加器实例,Executor端可以获得这个实例的副本,在处理后返回结果。最后,在Driver端可以对Executor最终的副本进行聚合。这与就实现了在不进行shuffle的时候也能够对多个Executor的数据进行分布式聚合。

spark为我们提供了默认的累加器,使我们可以进行简单的数据聚合。使用如下:
val rdd = sc.makeRDD(List(1, 2, 3, 4))
//获取系统累加器
//spark默认提供了简单数据聚合的累加器
val sumACC = sc.longAccumulator("sum")
//将num与累加器进行加和
rdd.foreach(sumACC.add(_))
//获取累加器的值
println(sumACC.value)
(3)累加器使用时的问题
若累加器在转换算子中调用,则会出现一些问题:
//当累加器出现在转换算子中的问题分析
//少加问题:不调用执行算子时,累加器不会执行
val mapRDD = rdd.map(sumACC.add(_))
println(sumACC.value) //0
//多加问题:行动算子多次调用转换算子,转换算子内的累加器会多次执行累加
mapRDD.collect()
mapRDD.collect()
println(sumACC.value) //20
由于累加器直接放在转换算子会出现 多加 & 少加 问题,因此我们一般会把累加器放在行动算子中,使其在最后阶段才执行累加操作。
(4)自定义累加器的实现
我们想要避免通过reduceByKey这类shuffle操作进行wordCount以节省时间,因此设计一个wordCount累加器,用于配合foreach进行wordCount操作。
a 自定义累加器的初步定义
我们可以参考spark内置的累加器定义:

可以看到,一个累加器,需要继承AccumulatorV2这个类。AccumulatorV2类存在泛型,分别是 累加器输入参数IN & 累加器返回参数OUT:

对于我们的wordCount累加器,很显然,输入参数IN为这个字符串本身,而输出参数OUT则为这个字符串 & 字符串出现的次数。因此,IN参数泛型为String,而OUT参数为mutable.Map[(String, Int)]。定义为可变Map是显然的,因为这个map需要随着foreach遍历不断变化。因此,我们自定义累加器初步定义如下:

b 在spark程序中初步注册并使用这个累加器

与使用内置累加器唯一的不同之处在于,由于内置累加器内部定义了注册方法,因此在创建时会自动向spark进行注册:

而我们自定义的累加器并未实现该方法,因此需要在创建后向spark进行注册,才能够去使用这个累加器进行分布式累加:

c 自定义累加器的累加方法实现

接下来我们只需要将自定义累加器中的抽象方法实现即可:
/**
* 1、继承AccumulatorV2,定义泛型:
* IN:累加器的输入数据类型,此处为wordCount,因此输入String类型
* OUT:累加器返回的数据类型,wordCount返回的是(单词, 出现次数),因此返回一个可变Map
*
*/
class MyAccumulator extends AccumulatorV2[String, mutable.Map[String, Long]] {
//用于进行统计累加结果的集合变量
private var wordCountMap = mutable.Map[String, Long]()
//判断当前累加器是否为初始状态
override def isZero: Boolean = wordCountMap.isEmpty
//复制累加器
override def copy(): AccumulatorV2[String, mutable.Map[String, Long]] = new MyAccumulator
//重置累加器,清空累加器的存储数据
override def reset(): Unit = wordCountMap.clear()
//当前Executor累加器的累加计算
override def add(word: String): Unit = {
//若不存在,则getOrElse获取的默认值为0,在该基础上 + 1
//若存在,getOrElse获取的值为之前的出现次数,在该基础上 + 1
wordCountMap.update(word, 1 + wordCountMap.getOrElse(word, 0))
}
//在Driver端合并多个Executor中累加器的方法
override def merge(other: AccumulatorV2[String, mutable.Map[String, Long]]): Unit = {
//获取当前累加器的wordCountMap
val thisMap = this.wordCountMap
//获取另一个累加器的wordCountMap
val otherWordCountMap = other.value
//合并两个Map,对当前的累加器进行值的修改(注意遍历 & 合并方向!!)
otherWordCountMap.foreach( //遍历otherWordCountMap中的所有元组
//对thisMap进行更新,将otherWordCountMap中的值更新到thisMap中
// 若thisMap中有该word,那么加上count
// 若thisMap中无该word,那么getOrElse默认为0,count = otherWordCountMap.count
{ case (word, count) => thisMap.update(word, count + thisMap.getOrElse(word, 0)) }
)
}
//累加器返回结果:将统计集合进行返回
override def value: mutable.Map[String, Long] = wordCountMap
}
以下是累加器方法解析:
① 定义当前累加器存储累加结果的属性

我们需要有一个变量,用于对每次累加后的结果进行记录,以便在下一次累加时使用该结果并对其进行更新。同时,这个变量也是当前累加器累加后进行Driver端合并的参数。
因为我们正在设计进行wordCount的累加器,所以每次的累加结果应该是wordCount期望的返回结果,也就是 (word, count) 的元组集合,所以定义为mutable.Map[String, Long]。
② add方法——当前Executor的累加器进行累加

由于我们是在进行wordCount,实际上就是在对这个Map[String, Long]进行更新。我们采用getOrElse的方式,无论 key = word 是否存在,都能够获得对应的value(key存在则获得对应value,key不存在则获得默认值0L),然后对value进行 + 1的计数。
注意,由于Map中元组类型是[String, Long],因此getOrElse的初始值需要是Long类型的,也就是需要设置为 0L 。
通过getOrElse获取value是scala中对Map集合进行更新的重要技巧。
③ value方法——将当前Executor的最终累加结果进行返回

没啥值得注意的,把最终的wordCountMap返回就行了。
④ merge方法——将多个Executor的累加器在Driver端进行合并

由于累加器的返回结果是Map[String, Long],因此该方法实际上就是在对两个Map进行合并。
注意事项:
- 在之前scala学习中,对两个Map合并可以采用foldLeft方法,但是在这里是有问题的。原因之一是:foldLeft并不是原地更新,会生成一个新的Map。那能否将更新结果赋值给当前累加器.map呢?也是不行的。因为foldLeft方法返回值为immutable的子类,与累加器的属性类型冲突。
- 我们使用foreach方法,遍历某一个Map,更新另一个Map。在进行这个合并操作时,需要特别注意合并顺序。由于我们要更新的是当前累加器的map,即thisMap,因此我们的目的是将otherWordCountMap中的所有元组更新到thisMap中。所以我们foreach遍历的应该是otherWordCountMap中的元组,才能保证该需求。
- foreach方法的偏函数匹配的是otherWordCountMap中的(word, count)元组,因此对thisMap进行更新时,调用getOrElse的应该是thisMap,这样才能进行thisMap中是否有otherWordCountMap元素的判断。
- 注意getOrElse的默认值类型,需要保持与Map类型一致,为Long,所以需要设置为 0L 。
⑤ isZero方法——判断累加器是否为初始状态

初始状态,即啥也没有,也就是wordCountMap这个结果集合为空。
⑥ copy方法——复制累加器

创建一个新的自定义累加器即可。
⑦ reset方法——重置累加器

重置,也就是清空。累加器的数据存储在wordCountMap集合中,因此清空该集合即可。
d 最终效果展示
object spark_acc04 {
def main(args: Array[String]): Unit = {
val sparkConf = new SparkConf().setMaster("local[*]").setAppName("acc")
val sc = new SparkContext(sparkConf)
val rdd = sc.makeRDD(List("scala", "spark", "world"))
//自定义累加器实现WordCount,避免使用reduce出现shuffle
val wordCountACC = new MyAccumulator
//向spark注册自定义累加器
sc.register(wordCountACC, "wordCountACC")
//通过foreach进行wordCount
rdd.foreach(wordCountACC.add(_))
//获取累加器的输出
println(wordCountACC.value)
sc.stop()
}
//自定义累加器
/**
* 1、继承AccumulatorV2,定义泛型:
* IN:累加器的输入数据类型,此处为wordCount,因此输入String类型
* OUT:累加器返回的数据类型,wordCount返回的是(单词, 出现次数),因此返回一个可变Map
*
*/
class MyAccumulator extends AccumulatorV2[String, mutable.Map[String, Long]] {
//用于进行统计累加结果的集合变量
private var wordCountMap = mutable.Map[String, Long]()
//判断当前累加器是否为初始状态
override def isZero: Boolean = wordCountMap.isEmpty
//复制累加器
override def copy(): AccumulatorV2[String, mutable.Map[String, Long]] = new MyAccumulator
//重置累加器,清空累加器的存储数据
override def reset(): Unit = wordCountMap.clear()
//当前Executor累加器的累加计算
override def add(word: String): Unit = {
//若不存在,则getOrElse获取的默认值为0,在该基础上 + 1
//若存在,getOrElse获取的值为之前的出现次数,在该基础上 + 1
wordCountMap.update(word, 1 + wordCountMap.getOrElse(word, 0L))
}
//在Driver端合并多个Executor中累加器的方法
override def merge(other: AccumulatorV2[String, mutable.Map[String, Long]]): Unit = {
//获取当前累加器的wordCountMap
val thisMap = this.wordCountMap
//获取另一个累加器的wordCountMap
val otherWordCountMap = other.value
//合并两个Map,对当前的累加器进行值的修改(注意遍历 & 合并方向!!)
otherWordCountMap.foreach( //遍历otherWordCountMap中的所有元组
//对thisMap进行更新,将otherWordCountMap中的值更新到thisMap中
// 若thisMap中有该word,那么加上count
// 若thisMap中无该word,那么getOrElse默认为0,count = otherWordCountMap.count
{ case (word, count) => thisMap.update(word, count + thisMap.getOrElse(word, 0L)) }
)
}
//累加器返回结果:将统计集合进行返回
override def value: mutable.Map[String, Long] = wordCountMap
}
}

5.3.2 广播变量
广播变量是分布式只读变量,用来高效分发较大的对象。向所有工作节点发送一个较大的只读值,以供一个或多个Spark操作使用。
比如,如果你的应用需要向所有节点发送一个较大的只读查询表,广播变量用起来都很顺手。在多个并行操作中使用同一个变量,但是 Spark会为每个任务分别发送。
(1)广播变量的作用
我们通过以下场景来分析广播变量的作用:

假设场景为:
我们想要对上面两个RDD的数据进行按key的连接,形成(String, (rdd1.value, rdd2.value))的数据。注意,不是按key的分组,而是将rdd1和rdd2进行连接操作!
- 最直接的方式是通过join操作,实现按key的连接:

但是这种操作是涉及笛卡尔积的,在数据量大的情况下进行笛卡尔积,数据增长就很恐怖了。这对shuffle的影响也是很大的。
- 将数据转换格式,map操作也是可以的,所以我们可以通过map操作去实现:

这个操作的核心是将rdd2的数据取出,形成一个map。在进行RDD.map操作时,只需要将rdd1的数据与已经准备好的数据进行转换即可,不涉及到rdd1与rdd2之间的shuffle,相对join操作而言,大大提升了性能。
但是该操作同样存在问题:当执行RDD.map操作的Executor分区较多,那么分配的Task也较多。同时,由于闭包数据map等是以Task进行分发的,那这就导致一个Executor中所有Task包含大量重复的数据,数据冗余,同时也占用了大量内存:

- 通过广播变量,将rdd2的数据放在Executor内存中,实现多个Task之间的共享:

既然我们不想要每个Task都分发一个闭包数据,那么为什么不把闭包数据放在Executor的内存中,让所有Task均从中读取数据呢?
我们将Executor内存中的闭包数据称之为广播变量,这些数据可以被多个Task共享,避免了每个Task都传输相同的闭包数据。
但是需要注意,由于数据是需要进行共享的,因此广播变量是只读,不可变的!
(2)广播变量的使用
设置广播变量:val bc = sc.broadcast(要广播的数据)
访问广播变量:bc.value
val rdd1 = sc.makeRDD(List(("a", 1), ("b", 2), ("c", 3)))
val rdd2 = sc.makeRDD(List(("c", 4), ("b", 5), ("a", 6)))
//通过map算子实现join效果,不会进行shuffle(提前将rdd2的数据取出)
val map = rdd2.collect().toMap
//设置广播变量,将map作为共享变量分发到Executor内存,让多个Task共享
val bc = sc.broadcast(map)
val mapRDD = rdd1.map(
{ case (word, count) => (word, (count, bc.value.getOrElse(word, 0))) }
)
println(mapRDD.collect().mkString(","))
5.4 SparkCore案例实操——用户行为数据分析
5.4.1 数据准备 & 数据说明
在之前的学习中,我们已经学习了Spark的基础编程方式,接下来,我们看看在实际的 工作中如何使用这些API实现具体的需求。这些需求是电商网站的真实需求,所以在实现 功能前,咱们必须先将数据准备好。

上面的数据图是从数据文件中截取的一部分内容,表示为电商网站的用户行为数据,主 要包含用户的4种行为:搜索,点击,下单,支付。
数据规则如下:
- 数据文件中每行数据采用下划线分隔数据
- 每一行数据表示用户的一次行为,这个行为只能是4种行为的一种
- 如果搜索关键字为null,表示数据不是搜索数据
- 如果点击的品类ID和产品ID为-1,表示数据不是点击数据
- 针对于下单行为,一次可以下单多个商品,所以品类ID和产品ID可以是多个,id之 间采用逗号分隔,如果本次不是下单行为,则数据(品类ID和产品ID)采用null表示
- 支付行为和下单行为类似
详细字段说明:
| 编号 | 字段名称 | 字段类型 | 字段含义 |
| 1 | date | String | 用户点击行为的日期 |
| 2 | user_id | Long | 用户的ID |
| 3 | session_id | String | Session 的 ID |
| 4 | page_id | Long | 某个页面的ID |
| 5 | action_time | String | 动作的时间点 |
| 6 | search_keyword | String | 用户搜索的关键词 |
| 7 | click_category_id | Long | 某一个商品品类的ID |
| 8 | click_product_id | Long | 某一个商品的ID |
| 9 | order_category_ids | String | 一次订单中所有品类的ID集合 |
| 10 | order_product_ids | String | 一次订单中所有商品的ID集合 |
| 11 | pay_category_ids | String | 一次支付中所有品类的ID集合 |
| 12 | pay_products | String | 一次支付中所有商品的ID集合 |
| 13 | city_id | Long | 城市 id |
样例类:
//用户访问动作表
//用户访问动作表
case class UserVisitAction(
date: String, //用户点击行为的日期
user_id: Long, //用户的ID
session_id: String, //Session 的 ID
page_id: Long, //某个页面的ID
action_time: String, //动作的时间点
search_keyword: String, //用户搜索的关键词
click_category_id: Long, //某一个商品品类的ID
click_product_id: Long, //某一个商品的ID
order_category_ids: String, //一次订单中所有品类的ID集合
order_product_ids: String, //一次订单中所有商品的ID集合
pay_category_ids: String, //一次支付中所有品类的ID集合
pay_product_ids: String, //一次支付中所有商品的ID集合
city_id: Long //城市 id
)
5.4.2 需求一:Top10热门商品品类统计
(1)需求说明


先按照点击数排名,靠前的就排名高;如果点击数相同,再比较下 单数;下单数再相同,就比较支付数。
(2)代码实现
最终结果展示:

a 版本一:分别求点击数、下单数、支付数,最后聚合排序
// TODO: Top10热门品类统计
val sparkConf = new SparkConf().setMaster("local[*]").setAppName("HotTop10Analysis")
val sc = new SparkContext(sparkConf)
//1、读取原始日志数据并且直接进行切分
val rawRDD = sc.textFile("datas/user_visit_action.txt").map(_.split("_"))
rawRDD.cache() //通过缓存将rawRDD进行存储
//2、统计品类的点击数量(品类ID,点击数量)
// 筛选出统计的数据:如果点击的品类ID和产品ID为-1,表示数据不是点击数据
// 即过滤掉切分后数组index = 6(品类ID) ,index = 7(产品ID) 均为-1的数据
val clickRDD = rawRDD
//过滤掉不是点击的数据
.filter(arr => arr(6) != "-1" && arr(7) != "-1")
//对点击数据转换,形成一个(品类ID, 1)的元组
.map(arr => (arr(6), 1))
//按照key进行统计,对value求和
.reduceByKey(_ + _)
//3、统计品类的下单数量(品类ID,下单数量)
// 同样,先筛选出不是下单的数据,即切分后数组index = 8 处为null的数据
val orderRDD = rawRDD
.filter(_(8) != null) //过滤掉不是下单的数据
//将下单的所有品类进一步切分,对每个品类ID都形成(品类ID, 1)的数据
.flatMap(arr => arr(8).split(",").map((_, 1)))
.reduceByKey(_ + _)
//4、统计品类的支付数量(品类ID,支付数量)
// 筛选标准为,切分后数组index = 10处为null的数据
val payRDD = rawRDD
.filter(_(10) != null)
//该步骤与下单统计同理
.flatMap(arr => arr(10).split(",").map((_, 1)))
.reduceByKey(_ + _)
//5、将品类进行排序,并且取前10名
// 点击数量 -> 下单数量 -> 支付数量
// 元组默认排序:先比较第一个、再比较第二个、、、以此类推
// 可以将2、3、4的数据整合,形成(品类ID, (点击数量, 下单数量, 支付数量)),然后按照_._2的tuple进行排序
//通过广播变量,将orderRDD结果 & payRDD结果进行广播
val orderMap = orderRDD.collect().toMap
val payMap = payRDD.collect().toMap
val bc = sc.broadcast( (orderMap, payMap) )
val resultRDD = clickRDD.map(
kv => {
val id = kv._1
val clickCount = kv._2
val orderCount = bc.value._1.getOrElse(id, 0)
val payCount = bc.value._2.getOrElse(id, 0)
(id, (clickCount, orderCount, payCount))
}
).sortBy(_._2, false)
val result = resultRDD.take(10)
//6、将结果采集到控制台打印出来
result.foreach(println)
sc.stop()
b 版本二:一次性统计点击数、下单数、支付数
在版本一中,我们对点击数、下单数、支付数分别进行reduceByKey的统计,这无疑进行了多次重复的shuffle,是可以避免的。
我们可以将数据格式进行转换,在一次RDD访问中,直接将点击数、下单数、支付数都转换成(品类ID, (点击数, 下单数, 支付数))的格式,然后对这个整体进行reduceByKey。
- 格式转换具体为:
当场景为点击数据时:由于下单数、支付数为0,因此变为(ID, (1, 0, 0))
当场景为下单数据时,由于点击数、支付数为0,因此变为(ID, (0, 1, 0))
当场景为支付数据时,由于点击数、下单数为0,因此变为(ID, (0, 0, 1))
- 需要注意:
由于单条下单数据、支付数据中包含多条品类ID,因此需要对该数据包含的所有品类ID都进行格式转换,会生成一个Array。因此在格式转换时,需要对这两类数据结果集合进行扁平化操作。
- 具体代码实现:
// TODO: Top10热门品类统计
val sparkConf = new SparkConf().setMaster("local[*]").setAppName("HotTop10Analysis")
val sc = new SparkContext(sparkConf)
//1、读取原始日志数据并且直接进行切分
val rawRDD = sc.textFile("datas/user_visit_action.txt").map(_.split("_"))
rawRDD.cache() //通过缓存将rawRDD进行存储
//一次性统计点击数、下单数、支付数
// 点击数直接转换为(ID, (1, 0, 0))
// 下单数直接转换为(ID, (0, 1, 0))
// 支付数直接转换为(ID, (0, 0, 1))
val flatRDD = rawRDD.flatMap( // 由于下单数、支付数统计涉及一条数据多个品类id,因此需要flat进行扁平化
datas => {
//如果是点击场合
if (datas(6) != "-1")
List((datas(6), (1, 0, 0))) //单数据源,需要转换为List进行扁平化
//如果是下单场合
else if (datas(8) != "null") {
val ids = datas(8).split(",") //品类ID的集合
ids.map((_, (0, 1, 0)))
}
//如果是支付场合,处理逻辑与下单场合类似
else if (datas(10) != "null") datas(10).split(",").map((_, (0, 0, 1)))
//其余场合,不在统计范围内,返回空集合
else Nil
})
//对数据进行聚合
val reduceRDD = flatRDD.reduceByKey((t1, t2) => (t1._1 + t2._1, t1._2 + t2._2, t1._3 + t2._3))
//排序并取top10
val resultRDD = reduceRDD.sortBy(_._2, false).take(10)
resultRDD.foreach(println)
c 版本三:对版本二的改进,使用ACC避免reduceByKey操作
在版本二中,对value的count仍然使用了reduceByKey的方法,虽然相比版本一少了很多shuffle,但是shuffle操作仍然是存在的。
为了避免shuffle操作,我们可以使用ACC,在对value的count时使用遍历 + 累加器,这样就可以避免使用reduceByKey了。
但是此操作难度较高,请记住:完成比完美更重要,优先完成,再考虑完美。
① 累加器的定义
/**
* 自定义累加器实现点击数、下单数、支付数的累加
* 输入参数IN:一个个数据处理后的元组,如:(品类ID, (0, 0, 1))
* 输出参数OUT:累加后的结果,即:(品类ID, (点击数, 下单数, 支付数))
*/
class HotTop10Accumulator extends AccumulatorV2[(String, (Int, Int, Int)), mutable.Map[String, (Int, Int, Int)]] {
private var hotTop10ResultMap = mutable.Map[String, (Int, Int, Int)]() //用于进行累加的结果集合
override def isZero: Boolean = hotTop10ResultMap.isEmpty
override def copy(): AccumulatorV2[(String, (Int, Int, Int)), mutable.Map[String, (Int, Int, Int)]] = new HotTop10Accumulator()
override def reset(): Unit = hotTop10ResultMap.clear()
//对累加器的累加操作
//对输入参数对应Key的value进行累加即可
override def add(v: (String, (Int, Int, Int))): Unit = {
val id = v._1
val values = v._2 //当前输入参数的value元组
val counts = hotTop10ResultMap.getOrElse(id, (0, 0, 0)) //当前累加器对应该id的累加结果
//将当前传入参数与结果集合进行累加
hotTop10ResultMap.update(id, (values._1 + counts._1, values._2 + counts._2, values._3 + counts._3))
}
//将分区累加器的结果进行合并
override def merge(other: AccumulatorV2[(String, (Int, Int, Int)), mutable.Map[String, (Int, Int, Int)]]): Unit = {
val otherMap = other.value //另一个累加器的结果集合
val thisMap = this.hotTop10ResultMap //当前累加器的结果集合
//遍历另一个累加器的结果集合,将所有结果更新到当前累加器的结果集合中
otherMap.foreach(
kv => {
val id = kv._1 //另一个累加器当前品类id
val otherCounts = kv._2 //另一个累加器当前品类id的统计结果
val thisCounts = thisMap.getOrElse(id, (0, 0, 0))
//将结果进行更新
thisMap.update(id, (thisCounts._1 + otherCounts._1, thisCounts._2 + otherCounts._2, thisCounts._3 + otherCounts._3))
}
)
}
//将累加器结果返回
override def value: mutable.Map[String, (Int, Int, Int)] = hotTop10ResultMap
}
- 关于输入/输出参数泛型:

这是累加器实现最重要的逻辑。
输入参数实际上就是我们对RDD遍历时RDD里面的数据格式。我们在RDD转换的时候,直接将所有类型的数据转换为(品类ID, (0/1, 0/1, 0/1))的格式,因此输入参数泛型即:(String, (Int, Int, Int))的元组。
输出参数也很好理解,我们要返回关于所有品类ID的相关信息统计,也就是说这个输出参数集合中,每一个数据都是(品类ID, (点击数, 下单数, 支付数))这样的格式,因此输出参数就是这样格式数据的集合。同时,由于输出参数需要进行累加,所以是mutable.Map[String, (Int, Int, Int)]格式。
- 关于累加器属性参数:
![]()
也比较简单,与输出参数泛型保持一致即可。
需要注意的是,我们只想让调用者通过Acc.value获取累加器结果,因此该参数使用了private进行私有化。
- 关于累加器add累加方法:

实际上就是将结果集合使用输入参数进行更新。
我们在方法中提前将输入参数的value元组取出;同时将结果集合对应key的value元组取出,若不存在,则赋予初始值(0, 0, 0)。这一步是为了便于更新操作。
由于输入参数一定是关于 点击/下单/支付 操作的任意一种,也就是说输入参数在某一位置是1,其余位置是0,因此进行更新时,只需要将values 和 counts在对应位置进行元素求和即可。
- 关于累加器merge合并累加器方法:

思路与add同理,提前取出两个累加器的value元组,然后进行对应位置元素的更新即可。
需要注意遍历方向是otherMap,因为要将另一个累加器的所有内容同步到当前累加器。所以,更新的是thisMap。
② 对累加器的使用
由于RDD中的数据是(品类ID, (0/1, 0/1, 0/1))的格式,因此此处直接展示如何使用这个累加器:
//通过累加器进行结果实现
//创建累加器对象
val hotTop10Accumulator = new HotTop10Accumulator
//注册累加器
sc.register(hotTop10Accumulator)
//对RDD内数据遍历,用累加器进行统计
flatRDD.foreach(hotTop10Accumulator.add)
//获取累加器结果,转为List,排序,取前10
val results = hotTop10Accumulator.value
③ 完整代码
def main(args: Array[String]): Unit = {
// TODO: Top10热门品类统计
val sparkConf = new SparkConf().setMaster("local[*]").setAppName("HotTop10Analysis")
val sc = new SparkContext(sparkConf)
//1、读取原始日志数据并且直接进行切分
val rawRDD = sc.textFile("datas/user_visit_action.txt").map(_.split("_"))
rawRDD.cache() //通过缓存将rawRDD进行存储
//一次性统计点击数、下单数、支付数
// 点击数直接转换为(ID, (1, 0, 0))
// 下单数直接转换为(ID, (0, 1, 0))
// 支付数直接转换为(ID, (0, 0, 1))
val flatRDD = rawRDD.flatMap( // 由于下单数、支付数统计涉及一条数据多个品类id,因此需要flat进行扁平化
datas => {
//如果是点击场合
if (datas(6) != "-1")
List((datas(6), (1, 0, 0))) //单数据源,需要转换为List进行扁平化
//如果是下单场合
else if (datas(8) != "null") {
val ids = datas(8).split(",") //品类ID的集合
ids.map((_, (0, 1, 0)))
}
//如果是支付场合,处理逻辑与下单场合类似
else if (datas(10) != "null") datas(10).split(",").map((_, (0, 0, 1)))
//其余场合,不在统计范围内,返回空集合
else Nil
})
//通过累加器进行结果实现
//创建累加器对象
val hotTop10Accumulator = new HotTop10Accumulator
//注册累加器
sc.register(hotTop10Accumulator)
//对RDD内数据遍历,用累加器进行统计
flatRDD.foreach(hotTop10Accumulator.add)
//获取累加器结果,转为List,排序,取前10
val results = hotTop10Accumulator.value.toList.sortBy(_._2).takeRight(10).reverse
results.foreach(println)
sc.stop()
}
/**
* 自定义累加器实现点击数、下单数、支付数的累加
* 输入参数IN:一个个数据处理后的元组,如:(品类ID, (0, 0, 1))
* 输出参数OUT:累加后的结果,即:(品类ID, (点击数, 下单数, 支付数))
*/
class HotTop10Accumulator extends AccumulatorV2[(String, (Int, Int, Int)), mutable.Map[String, (Int, Int, Int)]] {
private var hotTop10ResultMap = mutable.Map[String, (Int, Int, Int)]() //用于进行累加的结果集合
override def isZero: Boolean = hotTop10ResultMap.isEmpty
override def copy(): AccumulatorV2[(String, (Int, Int, Int)), mutable.Map[String, (Int, Int, Int)]] = new HotTop10Accumulator()
override def reset(): Unit = hotTop10ResultMap.clear()
//对累加器的累加操作
//对输入参数对应Key的value进行累加即可
override def add(v: (String, (Int, Int, Int))): Unit = {
val id = v._1
val values = v._2 //当前输入参数的value元组
val counts = hotTop10ResultMap.getOrElse(id, (0, 0, 0)) //当前累加器对应该id的累加结果
//将当前传入参数与结果集合进行累加
hotTop10ResultMap.update(id, (values._1 + counts._1, values._2 + counts._2, values._3 + counts._3))
}
//将分区累加器的结果进行合并
override def merge(other: AccumulatorV2[(String, (Int, Int, Int)), mutable.Map[String, (Int, Int, Int)]]): Unit = {
val otherMap = other.value //另一个累加器的结果集合
val thisMap = this.hotTop10ResultMap //当前累加器的结果集合
//遍历另一个累加器的结果集合,将所有结果更新到当前累加器的结果集合中
otherMap.foreach(
kv => {
val id = kv._1 //另一个累加器当前品类id
val otherCounts = kv._2 //另一个累加器当前品类id的统计结果
val thisCounts = thisMap.getOrElse(id, (0, 0, 0))
//将结果进行更新
thisMap.update(id, (thisCounts._1 + otherCounts._1, thisCounts._2 + otherCounts._2, thisCounts._3 + otherCounts._3))
}
)
}
//将累加器结果返回
override def value: mutable.Map[String, (Int, Int, Int)] = hotTop10ResultMap
}
5.4.3 需求二:Top10热门品类中每个品类的Top10活跃Session统计
(1)需求说明

也就是说,我们要选出热门品类的点击数据,统计每个品类中每个session出现的次数,然后对每个品类中session出现次数的前10进行选取。
(2)代码实现
由于是基于需求一的统计,因此关于热门品类ID的前十获取不再此展示。我们在这进行了需求一结果的处理,因为我们只需要品类ID,因此我们将需求一结果转换为了品类ID的集合。
//排序并取top10的品类ID集合
val ids = reduceRDD.sortBy(_._2, false).take(10).map(_._1)
println("热门品类ID为:" + ids.mkString(","))
//对相关session进行统计,统计点击了热门品类的所有session
//过滤非点击数据,并形成((品类ID,session), 1)的格式,便于按key进行统计求和
val sessionRDD = rawRDD
.filter(_(6) != "-1") //筛选点击数据
.map(datas => ((datas(6), datas(2)), 1)) //转换格式,方便按key求和
//判断上述sessionRDD中的品类ID是否位于热门品类ID集合中,筛选出在其中的数据
val hot10SessionRawRDD = sessionRDD
.filter(kv => ids.contains(kv._1._1)) //筛选出热门品类id的数据
//按key = (品类ID,session) 进行统计
val sumRDD = hot10SessionRawRDD
.reduceByKey(_ + _) //按key进行热门品类,session的点击数求和
//将数据转换为(品类ID, (session, sum))的格式,便于按品类ID分组
val groupRDD = sumRDD
.map(kv => (kv._1._1, (kv._1._2, kv._2))) //转换格式,使品类ID作为唯一key
.groupByKey() //按品类ID分组
//对每个品类ID的session统计按照sum降序排列,选出每个品类的点击数前10session
val resultRDD = groupRDD
.mapValues(_.toList.sortWith(_._2 > _._2).take(10))
val results = resultRDD.collect()
results.foreach(println)
5.4.4 需求三:页面单跳转换率统计
(1)需求说明

假设我们需要分析首页 -> 详情页面 的跳转率,那么我们就需要多少用户点击了首页,然后有多少用户从首页跳转到了详情页面。
这种页面跳转率是跟用户有关的,只有一个用户在一次点击中,从首页点击到详情页面才是一个有效的数据;若两个用户,一个点击了首页,另一个点击详情,这并不能被称为页面跳转。
同时,只有当一个用户在两次相邻(时间上有序)的点击中进行了首页 & 详情 的点击才算一次有效跳转。
我们以下列数据进行分析:
- 数据初步处理1:

首先,Colum1为所有用户的点击情况。由于我们要统计页面跳转情况,这是跟用户有关的,因此我们首先需要按照用户进行分组,所以就形成了Colum2。
同时,跳转情况是有序的,与点击时间是相关的。比如,对同一个用户,12.05首页 -> 12.07详情的跳转就是正确的,而12.08下单 -> 12.06详情的页面跳转就是错误的,因为这不是相邻的点击。所以,我们需要在按照用户分组的情况下,按照时间再进行一次排序,形成了Colum3。
- 数据预处理2:

由于我们已经获取了所有用户页面的点击情况(按时间排序),所以此时用户相关信息就已经不重要了,我们只需要将分组内数据保留页面数据即可。
基于上述的要求,我们获得了上图的数据,这是按照用户分组后,每个用户的页面点击情况。
- 分母计算:

分母计算其实很简单,我们只需要统计每个页面出现的次数就可以了。因为单纯分母不涉及有效的页面跳转,因此不需要按用户的点击状态去统计。
- 分子计算:

分子需要的数据是页面跳转情况,也就是对一个用户来说,按照时间,两个页面点击的情况。所以,我们在每个用户分组后的数据中,按时间两两组合,形成一个页面点击情况。比如对于第一个用户,他按照时间分布点击了首页、详情、下单、支付,那么他所有的页面跳转情况则为:首页 -> 详情、详情 -> 下单、 下单 -> 支付。其余同理。我们按照这样的处理,获取以下格式数据:

最后,将相同的页面跳转情况进行汇总,就可以得出所有的页面跳转情况了。
- 最终页面单跳率计算:
有了每个页面的被点击情况 & 页面跳转情况后,我们就可以计算最终页面单跳率了。比如计算首页 -> 详情的单跳率,就是多少用户在首页的情况下跳转到了详情,即:首页 -> 详情的分子汇总 / 首页的分母汇总。
(2)代码实现
val sparkConf = new SparkConf().setMaster("local[*]").setAppName("PageJumpRateAnalysis")
val sc = new SparkContext(sparkConf)
//读取数据,获得(sessionID, 页面ID, 动作时间)的初始数据
val rawRDD = sc.textFile("datas/user_visit_action.txt")
.map(_.split("_")).map(lines => (lines(2), lines(3), lines(4)))
rawRDD.cache()
/*
分母获取:每个页面出现的次数
*/
//前一步RDD中数据格式:(sessionID, 页面ID, 动作时间)
val pageCounts = rawRDD.map(tp => (tp._2, 1)).reduceByKey(_ + _)
//最终形成数据(页面ID, 出现次数),将结果保存到内存中,形成数据集
val pageCountsResults = pageCounts.collect().toMap
/*
分子获取:统计用户的页面跳转情况
*/
//按照sessionID分组
//前一步RDD中数据格式:(sessionID, 页面ID, 动作时间)
val groupRDD = rawRDD.groupBy(_._1)
//对每个分组的数据按照动作时间进行排序
//前一步RDD中数据格式:(sessionID, CompactBuffer((页面ID, 动作时间)))
val sortRDD = groupRDD.map(
group => {
//形成(session,List(页面ID))的数据集合,其中页面ID按照动作时间进行了排序
(
group._1, //sessionID
group._2.toList
.sortWith(_._3 < _._3) //按照动作时间排序
.map(_._2) //选出页面ID作为唯一元素
)
}
)
//统计每个用户的页面跳转情况,也就是对List(页面ID)进行进一步转换,形成(页面ID1, 页面ID2)的新集合
//前一步RDD中数据格式:(sessionID, List(页面ID))
val zipRDD = sortRDD.mapValues(list => {
//将list中第一个元素去除,形成一个新list
//将list与新list进行zip操作
val newList = list.takeRight(list.size - 1) //取出后n个元素
list.zip(newList)
})
//将所有用户的页面跳转情况进行汇总,按照相同的页面跳转情况进行分类并且统计出现次数
//前一步RDD中数据格式:(sessionID, List((页面1, 页面2)))
//首先去除session信息,将所有页面跳转情况集合扁平化
//然后将页面跳转情况转换为(页面跳转情况, 1)的数据格式
//最后按key统计每个页面跳转情况的综合
val pageJumpCountsRDD = zipRDD.flatMap(_._2).map((_, 1)).reduceByKey(_ + _)
//将结果保存到内存中,形成数据集
val pageJumpCountsResults = pageJumpCountsRDD.collect().toMap
/*
最终每个页面跳转率的计算
*/
//即:当前页面到别的页面的概率
//比如想要计算:首页 -> 详情页面 的跳转率,则使用:(首页, 详情页面)跳转次数 / 首页次数
//即:(A, B).counts / A.counts
val results = pageJumpCountsResults.map(
{ case ((pageId1, pageId2), sum) =>
val pageId1Counts = pageCountsResults.getOrElse(pageId1, 0)
((pageId1, pageId2), sum.toDouble / pageId1Counts)
}
)
results.foreach(datas => println(s"页面${datas._1._1}到页面${datas._1._2}的单跳转换率为: ${datas._2}"))
sc.stop()
a 数据预处理

b 分母获取(每个页面出现次数)

按页面ID分组,然后按key进行求和即可。
c 分子获取(每个页面跳转的出现次数)

- 首先按照sessionID分组,形成每个sessionID的页面点击情况,同时按照页面点击时间进行排序:

- 然后将每个session的相关相邻页面动作两两组合(在基于时间排序后的情况下),形成页面跳转情况:

在动作基于时间升序排列后,组合实际上是(index = n, index = n + 1)。此处通过zip方法进行组合,将原先的List取出第一个元素,那么新的List就是(index = 1,…index = n)的集合,而原先List就是(index = 0,…index = n)的集合,通过原先List.zip新List,就可以形成List((index = 0, index = 1),…,(index = n – 1, index = n))的页面跳转集合。
- 最后,去除sessionID相关信息,将所有的页面跳转情况扁平化,统计每个页面跳转情况出现次数(wordcount):

d 页面跳转率的计算

5.5 工程化代码
为了不同功能代码之间的分层解耦,在传统web开发中采用MVC架构。其中,model用来存放数据模型,View用来生成可视化效果,Controller用来控制model与view之间的传输。

关于数据服务相关架构(后端架构),采用Controller + Service + Dao,对数据处理进行进一步解耦。

5.5.1 对WordCount程序进行三层架构化
(1)一个工程的目录结构

- application:应用类,所有程序的启动入口
- bean:数据相关的实体类
- common:共同具有的属性
- controller:控制层相关代码
- dao:持久层相关代码
- service:服务层相关代码
- util:工具类相关代码
(2)Application代码

Application层为应用入口,功能很简单:设置spark环境 & 调用Controller(传递Controller所需要的参数)。
(3)Controller层代码

调用Service层代码,传递路径参数,让Service层执行计算逻辑。
(4)Service层代码

真正执行WordCount计算逻辑的代码。需要从Dao中获取数据源。
(5)Dao层代码

从path中通过textFile读取文件,生成原始RDD。sc是导入Application中的环境的变量。
5.5.2 架构优化——控制抽象
(1)通过控制抽象封装Application启动类特质
为了避免每个应用程序Application都要写对应的spark环境配置,我们可以将其抽取为共同的特质,通过混入的方式在每个Application中进行使用。
特质类的抽取代码为:

- 由于master & appName两个参数可以由不同应用传入,所以将二者设置为形参,并赋予其初始值,在不传入的时候使用默认值进行环境配置。
- 由于每个应用的操作不同,因此我们采用控制抽象,让调用者直接将操作代码块在调用方法时传入。由于这些操作需要在环境配置后传入,所以采用柯里化,在第一个返回函数的调用中添加op代码块,作为此次应用要执行的任务代码。
具体启动类混入启动特质并使用:

由于我们要调用Controller进行任务调度,所以op参数就是wordController.dispatch()的代码块。
(2)Controller & Service层的特质封装
同样的,每个任务都会有Controller & Service的调用,我们可以也将其抽取成特质,放到common文件夹,在具体的Controller & Service中混入调用即可。
该过程可以抽象为Service & ServiceImpl的解耦。
Controller的抽取 & 混入使用:
- 特质中定义抽象方法dispatch,作为所有Controller的调度方法模板:

- 在具体Controller中重写该方法,赋予具体逻辑:

Service的抽取 & 混入使用:
- 在特质中定义抽象方法dataAnalysis,作为所有Service数据分析的模板:

- 在具体Service中重写该抽象方法,赋予具体的数据分析逻辑:

(3)Dao层特质封装——SparkContext的获取
由于Dao层对数据的读取方法基本一致,所以可以直接在特质中定义具体的文件读取方法:

但是我们会发现,特质中无法获取SparkContext对象,因此sc是爆红的。那么现在首要问题就是如何获取SparkContext到Dao的环境中。
我们首先需要知道,我们对Controller、Service、Dao的三层解耦是对代码的解耦,但是一个sparkJob是由三者共同完成的,因此在同一个sparkJob中三者位于同一个线程中。那么,我们何不将信息存放在线程中,由需要使用的代码(Dao)进行读取呢:

Scala是可以使用Java的类库的,而Java中提供了ThreadLocal类,让我们可以将信息存储到ThreadLocal中。
首先我们创建一个util工具类,专门用于创建ThreadLocal对象并且将需要的环境内容放到ThreadLocal中:

ThreadLocal已经为我们提供了设置对象(set)取出对象(get)、移除对象(remove)方法,我们只需要对其进行进一步封装即可。
此处需要注意,ThreadLocal的生命周期是在一个线程运行时,若线程结束则不再存在。也就是说,关于ThreadLocal的使用需要在同一个线程内的方法完成,不能放在方法外部,否则将只能访问到空对象。
然后我们就可以在TDao层进行SparkContext的获取并使用了:

更多推荐
所有评论(0)