本章为Spark Core的详细内容

Spark3.x指北全系列目录:

Spark基础概念请看:Spark3.x指北——1:Spark基础概念

SparkSQL内容请看:Spark3.x指北——3:SparkSQL

SparkStreaming内容请看:Spark3.x指北——4:SparkStreaming

目录

5 Spark Core

5.1 引入——通过网络编程模拟分布式计算

5.1.1 基础的客户端向服务端发送数据演示

(1)我们首先需要一个服务器,用来等待客户端连接、接收数据、输出数据:

(2)然后我们需要一个客户端,向服务器发送数据:

5.1.2 将客户端发送的数据变为发送一个计算任务(数据 & 操作硬编程)

(1)现在我们需要额外添加一个计算任务类,用于存储需要计算的数据 & 需要进行的计算操作:

(2)然后我们将客户端发送内容进行更改,通过ObjectOutputStream发送计算任务对象:

(3)最后,我们在服务器通过ObjectInputStream接收这个计算任务,然后执行这个任务里的计算方法,得出计算结果即可:

5.1.3 将计算任务进行分布式处理

(1)我们重新定义一个任务类,这个任务类的 计算数据 和 计算操作 均不赋值,只是封装一个compute方法,这个方法中对 数据 进行 操作。我们可以将这个类理解为控制抽象,仅提供对数据的规范和对操作的规范,具体的数据和操作均由调用者提供:

(2)保留5.1.2的提交任务类,这个类将会为Task提供数据和操作:

(3)我们在客户端将SubTask切分为两部分,分别赋值给Task1和Task2,然后将这两个任务发送给两台服务器9999和9998:

(4)服务器端分别设置两个端口接收客户端发送的两个任务,分别执行这两个任务然后输出结果:

(5)对于5.1.3分布式计算模拟的总结

5.2 数据结构——RDD

5.2.1 什么是RDD

5.2.2 从IO的装饰者模式到RDD的组合

5.2.3 RDD的五大特性

(1)分区列表

(2)分区计算函数

(3)RDD依赖关系的获取(血统)

(4)分区器(KV数据类型时采用的分区策略)

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

5.2.4 RDD的执行原理

(1)启动Yarn集群,获取资源调度框架,用于管理资源:

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

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

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

5.2.5 RDD基础编程——创建RDD

(1)从内存中直接创建RDD——以内存中数据为数据源

a 通过 parallelize 方法

b 通过 makeRDD 方法创建(常用)

(2)从外部存储中创建RDD

a 通过textFile方法

Ⅰ 路径参数默认以当前环境根目录问基准,绝对路径 / 相对路径均可

Ⅱ 当我只写目录名时,我就可以读取目录中所有文件

Ⅲ 可以使用通配符,筛选目录中的文件

Ⅳ 路径不仅可以为本地文件,同时也可以是HDFS、HBase中的数据

b 通过wholeTextFiles方法创建

(3) 从其他RDD创建

(4) 直接创建RDD(new)

5.2.6 RDD基础编程——分区与并行度

(1)从内存中读取数据时创建分区

a 分区的设定

b 分区数据的分配

(2)从文件中读取数据时创建分区

5.2.7 RDD基础编程——转换算子

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

a 基础使用——将RDD数据集中每一个数据都执行 * 2 的转换

b 一个小实操案例

c 并行计算效果演示

(2)mapPartitions算子——对分区进行整体处理

a 基础使用——一次性对一个分区的数据进行 * 2 的转换操作

b 案例演示

c 两个算子:map & mapPartitions の对比

(3)mapPartitionsWithIndex算子

a 通过该方法获取第二个分区的所有数据

b 通过该方法获取每个数据所在的分区索引

(4)flatMap算子——扁平化

a 将RDD中两个数据集合进行扁平化合并

b 将RDD中的字符串集合元素按空格拆分并且扁平化

c 对于a的进阶,将混合类型的数据集合扁平化

(5)glom算子——分区数据集类型的转换:List => Array

a 基础使用

b 将分区内的最大值取出,进行分区间最大值求和

c 分区不变的含义

(6)groupBy算子——分布式分组

a 基础使用

b shuffle概念了解

c 统计apache.log中每个时间段访问数据量

(7)filter算子——过滤

a 基础使用

b 数据倾斜

c 选出apache.log中 2015年5月17日的请求路径

(8)sample算子——抽样

(9)distinct算子——分布式去重

a 基础使用

b 去重原理分析(分布式计算的去重计算)

(10)coalesce算子——改变分区数量(默认用于缩减分区)

a 基础使用——缩减分区

b 数据倾斜 & shuffle

c 使用coalesce实现增加分区

(11)repartition算子——增加分区

(12)sortBy算子——根据规则函数对RDD数据进行排序

(13)针对两个RDDs数据的操作——交 & 并 & 差 & 拉链

a 并集——union

b 交集——intersection

c 差集——subtract

d 拉链——zip

e 关于双value操作的注意事项

(14)针对KV类型数据的操作

a partitionBy算子——按照分区规则重新分区

① partitionBy的相关说明

② partitonBy的基础使用

③ partitionBy的注意事项

b reduceByKey算子——按照Key进行聚合

① 基础使用

② 注意事项

c groupByKey算子——按照Key进行分组

d 区分:reduceByKey & groupByKey的区别和使用场景

e aggregateByKey算子——对分区内和分区间分别指定计算规则

① 形参说明

② 实操——分区内按Key求最大值,分区间按Key进行求和

③ 对aggregateByKey的深入理解 & 案例实现

f foldByKey算子——分区间和分区内计算规则相同的简化

g combineByKey算子——直接将第一个数据转换作为初始值进行计算

h 区分:reduceByKey & aggregateByKey & foldByKey & combineByKey的区别以及使用场景

i join & leftOuterJoin & rightOuterJoin算子——连接操作

① join算子

② leftOuterJoin & rightOuterJoin算子

j cogroup算子——不同RDD内相同key的分别聚合

(15)转换算子案例实操

a 数据源部分数据如下:

b 实现思路:

c 代码实现:

5.2.8 RDD基础编程——行动算子

(1)reduce算子——将RDD数据聚合并返回

(2)collect算子——将RDD数据收集到内存形成结果数组

(3)count算子——计算RDD中的数据个数

(4)first算子——取出RDD中第一个数据

(5)take算子——取出RDD中前n个数据形成Array

(6)takeOrdered算子——返回RDD中按照排序规则排序后的数据的前n个

(7)aggregate算子——分别与初始值进行分区内 & 分区间运算

(8)fold算子——aggregate算子的简化,分区内 & 分区间计算规则相同

(9)countByValue & countByKey算子——统计value出现的次数& 统计Key出现的次数

a countByValue算子

b countByKey算子

(10)save相关算子——保存结果到文件中

a saveAsTextFile——保存为text文件:

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

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

d 使用展示

(11)foreach算子——分布式遍历

5.2.9 序列化

(1)引入——从foreach算子遍历对象到序列化

(2)RDD的闭包检测(Driver端执行)

(3)Kryo序列化(大数据场景下序列化的改善)

5.2.10 依赖关系

(1)血缘关系

(2)宽窄依赖

a 窄依赖(OneToOne依赖)

b 宽依赖(Shuffle依赖)

(3)RDD阶段划分

a 概念解析

① 窄依赖时的阶段划分

② 宽依赖时的阶段划分

b 源码分析

① 阶段划分源码的调用情况

② 阶段划分的具体执行

③ 阶段划分的总结

(4)RDD任务划分

a 任务划分

b 任务划分的源码解读

5.2.11 持久化

(1)持久化的引入

(2)在代码中进行持久化

(3)持久化的作用总结

(4)CheckPoint检查点

(5)持久化 & CheckPoint 的区别

5.2.12 分区器

(1)Hash分区器

(2)Range分区器

(3)自定义分区器的实现案例

5.2.13 文件的读取 & 保存(保存具体源码参考5.2.8-save相关算子)

(1)text文件读取 & 保存

(2)sequence文件读取 & 保存

(3)object文件读取 & 保存

5.3 数据结构——累加器 & 广播变量

5.3.1 累加器(不用shuffle的分布式聚合)

(1)引入——没有累加器时分布式计算的问题

(2)累加器的原理 & 简单使用

(3)累加器使用时的问题

(4)自定义累加器的实现

a 自定义累加器的初步定义

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

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

① 定义当前累加器存储累加结果的属性

② add方法——当前Executor的累加器进行累加

③ value方法——将当前Executor的最终累加结果进行返回

④ merge方法——将多个Executor的累加器在Driver端进行合并

⑤ isZero方法——判断累加器是否为初始状态

⑥ copy方法——复制累加器

⑦ reset方法——重置累加器

d 最终效果展示

5.3.2 广播变量

(1)广播变量的作用

(2)广播变量的使用

5.4 SparkCore案例实操——用户行为数据分析

5.4.1 数据准备 & 数据说明

5.4.2 需求一:Top10热门商品品类统计

(1)需求说明

(2)代码实现

a 版本一:分别求点击数、下单数、支付数,最后聚合排序

b 版本二:一次性统计点击数、下单数、支付数

c 版本三:对版本二的改进,使用ACC避免reduceByKey操作

① 累加器的定义

② 对累加器的使用

③ 完整代码

5.4.3 需求二:Top10热门品类中每个品类的Top10活跃Session统计

(1)需求说明

(2)代码实现

5.4.4 需求三:页面单跳转换率统计

(1)需求说明

(2)代码实现

a 数据预处理

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

c 分子获取(每个页面跳转的出现次数)

d 页面跳转率的计算

5.5 工程化代码

5.5.1 对WordCount程序进行三层架构化

(1)一个工程的目录结构

(2)Application代码

(3)Controller层代码

(4)Service层代码

(5)Dao层代码

5.5.2 架构优化——控制抽象

(1)通过控制抽象封装Application启动类特质

(2)Controller & Service层的特质封装

(3)Dao层特质封装——SparkContext的获取


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。我们开始迭代:

  1. 第一次迭代:i = 0,start = (i * length) / numSlices = (0 * 5) / 3 = 0,end = ((i + 1) * length) / numSlices = ((0 + 1) * 5) / 3 = 1,因此第一个范围为(0, 1)。
  2. 第二次迭代:i = 1,则start = (1 * 5) / 3 = 1,end = ((1 + 1) * 5) / 3 = 3,第二个范围为(1, 3)。
  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 实现思路:
  1. 获取原始数据,将其按空格拆分,形成 (时间戳, 省份, 城市, 广告, 用户) 类型的数据。
  2. 将拆分后数据选出 (省份,广告) ,并且对每一条数据都map成 ((省份, 广告), 1) 的形式。这个1表示当前广告的出现次数,我们对每一条数据都采用1作为默认值,便于我们后续进行求和统计。
  3. 按照第二步中的数据形式,按key对value进行求和,统计这个广告在这个省内出现的次数,将数据map成 ((省份, 广告), 出现次数)
  4. 将统计求和后的数据进行转换,map成 (省份, (广告, 出现次数)) 的形式,让省份成为唯一key,便于后续基于省份划分广告。
  5. 对第三步中的数据按省份进行分组,并且进行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的获取并使用了

更多推荐