一、转换算子
1.map
package sparkCore.rddTransform
import org.apache.spark.sql.SparkSession
object RDDTransformV1 {
def main(args: Array[String]): Unit = {
val spark = SparkSession
.builder()
.appName(this.getClass.getName)
.master("local[*]")
.getOrCreate()
val sc = spark.sparkContext
val rdd01 = sc.makeRDD(List(1,2,3,4))
val rdd01A = rdd01.map(func01)
val rdd01B = rdd01.map(x => x*10) //匿名函数方式,x类型可推导,无需再写类型
//分布式执行测试
val rdd02 = sc.makeRDD(List(1,2,3,4,5),2)
val rdd02A = rdd02.map {
x => {
println(s"aaaaaaaaaaaaa:$x")
x*10
}
}
val rdd02B = rdd02A.map {
x => {
println(s"bbbbbbbbbbbbb:$x")
x
}
}
rdd02B.collect()
sc.stop()
}
val func01 = (x: Int) => x*10
}
2.mapPartitions
package sparkCore.rddTransform
import org.apache.spark.sql.SparkSession
object RDDTransformV2 {
def main(args: Array[String]): Unit = {
val spark = SparkSession
.builder()
.appName(this.getClass.getName)
.master("local[*]")
.getOrCreate()
val sc = spark.sparkContext
val rdd01 = sc.makeRDD(List(1,2,3,4),4)
//以分区为单位进行批量处理
val rdd01A = rdd01.mapPartitions {
iter => {
println("*****************") //多少个分区执行多少次
iter.map(_ * 10)
}
}
val rdd02 = sc.makeRDD(List(1,2,3,4,5,6,7,8),4)
//获取每个分区最大值
val rdd02A = rdd02.mapPartitions {
// iter => iter.toList.sortWith(_ > _).take(1).iterator
iter => List(iter.max).iterator
}
val rdd02B = rdd02A.mapPartitionsWithIndex {
(index, iter) => {
iter.foreach {
max_val => println(s"index:$index max_val:$max_val")
}
iter
}
}
rdd02B.collect()
sc.stop()
}
}
3.mapPartitionsWithIndex
package sparkCore.rddTransform
import org.apache.spark.sql.SparkSession
object RDDTransformV3 {
def main(args: Array[String]): Unit = {
val spark = SparkSession
.builder()
.appName(this.getClass.getName)
.master("local[*]")
.getOrCreate()
val sc = spark.sparkContext
val rdd01 = sc.makeRDD(List(1, 2, 3, 4, 5, 6, 7), 4)
//操作指定分区的数据
val rdd01A = rdd01.mapPartitionsWithIndex {
(index, iter) => {
if (index % 2 == 1) {
iter.map((_, s"${index}_A"))
} else {
//Nil.iterator
iter.map((_, s"${index}_B"))
}
}
}
println(rdd01A.collect.toBuffer)
sc.stop()
}
}
4.flatMap
package sparkCore.rddTransform
import org.apache.spark.rdd.RDD
import org.apache.spark.sql.SparkSession
object RDDTransformV4 {
def main(args: Array[String]): Unit = {
val spark = SparkSession
.builder()
.appName(this.getClass.getName)
.master("local[*]")
.getOrCreate()
val sc = spark.sparkContext
val rdd01: RDD[List[Int]] = sc.makeRDD[List[Int]](List(List(1,2), List(3,4), List(5,6)),2)
val rdd02: RDD[String]= sc.makeRDD[String](List("hello scala", "hello spark","hello java"),2)
val rdd03: RDD[Any] = sc.makeRDD(List(List(1,2), 3.14, "hello spark"),2)
//传入的函数需要返回一个可迭代的集合
val rdd01A: RDD[Int] = rdd01.flatMap(list => list.filter(_ % 2 == 0))
val rdd02A: RDD[String] = rdd02.flatMap(word => word.split(" "))
//使用模式匹配去处理不同类型数据,也可以传入一个偏函数去处理
val rdd03A = rdd03.flatMap {
elem => {
elem match {
case list: List[Int] => list
case num: Double => List(num)
case word: String => word.split(" ")
case _ => Nil
}
}
}
//模式匹配可转换成偏函数
val rdd03B = rdd03.flatMap(partialFunc)
println(rdd03B.collect.toBuffer)
sc.stop()
}
def partialFunc: PartialFunction[Any,List[Any]] = {
case list: List[Int] => list
case num: Double => List(num)
case word: String => word.split(" ").toList
case _ => Nil
}
}
5.groupBy
package sparkCore.rddTransform
import org.apache.spark.sql.SparkSession
object RDDTransformV5 {
def main(args: Array[String]): Unit = {
val spark = SparkSession
.builder()
.appName(this.getClass.getName)
.master("local[*]")
.getOrCreate()
val sc = spark.sparkContext
val rdd01 = sc.makeRDD(List(1,2,3,4,5,6,7),4)
//将数据源中每一个数据进行分组判断,根据返回的key进行分组
val rdd01A = rdd01.groupBy(groupFunc01)
val rdd01B = rdd01.groupBy((num: Int) => num % 2) //匿名函数
val rdd01C = rdd01.groupBy(_ % 2)
//println(rdd01A.collect.toBuffer) //ArrayBuffer((0,CompactBuffer(2, 4, 6)), (1,CompactBuffer(1, 3, 5, 7)))
val rdd02 = sc.makeRDD(List("华为_A","华为_B","苹果_C","苹果_D","小米_E","小米_F"),4)
val rdd02A = rdd02.groupBy(groupFunc02)
val rdd02B = rdd02.groupBy((phone: String) => phone.split("_")(0))
val rdd02C = rdd02.groupBy(_.split("_")(0))
println(rdd02C.collect.toBuffer) //ArrayBuffer((华为,CompactBuffer(华为_A, 华为_B)), (小米,CompactBuffer(小米_E, 小米_F)), (苹果,CompactBuffer(苹果_C, 苹果_D)))
sc.stop()
}
val groupFunc01: Int => Int = num => {
num % 2
}
def groupFunc02(phone: String): String = {
phone.split("_")(0)
}
}
6.sample
package sparkCore.rddTransform
import org.apache.spark.sql.SparkSession
object RDDTransformV6 {
def main(args: Array[String]): Unit = {
val spark = SparkSession
.builder()
.appName(this.getClass.getName)
.master("local[*]")
.getOrCreate()
val sc = spark.sparkContext
val rdd01 = sc.makeRDD(List(1,2,3,4,5,6,7,8,9,10))
val rdd01A = rdd01.sample(
false,
0.5,
1
)
println(rdd01A.collect.mkString(","))
sc.stop()
}
}
7.distinct
package sparkCore.rddTransform
import org.apache.spark.sql.SparkSession
object RDDTransformV7 {
def main(args: Array[String]): Unit = {
val spark = SparkSession
.builder()
.appName(this.getClass.getName)
.master("local[*]")
.getOrCreate()
val sc = spark.sparkContext
val rdd01 = sc.makeRDD(List(1,1,2,2,3,3))
val rdd01A = rdd01.distinct()
println(rdd01A.collect.mkString(","))
sc.stop()
}
}
8.coalesce、repartition
package sparkCore.rddTransform
import org.apache.spark.sql.SparkSession
object RDDTransformV8 {
def main(args: Array[String]): Unit = {
val spark = SparkSession
.builder()
.appName(this.getClass.getName)
.master("local[*]")
.getOrCreate()
val sc = spark.sparkContext
val rdd01 = sc.makeRDD(List(1,2,3,4,5,6,7,8),3)
val rdd01A = rdd01.coalesce(2)
val rdd01B = rdd01.coalesce(
5,
true
)
val rdd01C = rdd01.repartition(5)
println(rdd01C.collect.mkString(","))
sc.stop()
}
}
9.sortBy
package sparkCore.rddTransform
import org.apache.spark.sql.SparkSession
object RDDTransformV9 {
def main(args: Array[String]): Unit = {
val spark = SparkSession
.builder()
.appName(this.getClass.getName)
.master("local[*]")
.getOrCreate()
val sc = spark.sparkContext
val rdd01 = sc.makeRDD(List("A:1", "B:3", "C:5", "D:4", "E:2", "F:0"), 3)
val rdd01A = rdd01.sortBy(str => str.split(":")(1), false)
val rdd01B = rdd01.sortBy(_.split(":")(1),false)
println(rdd01B.collect.mkString("n"))
sc.stop()
}
}
10.intersection、union、subtract、zip
package sparkCore.rddTransform
import org.apache.spark.sql.SparkSession
object RDDTransformV10 {
def main(args: Array[String]): Unit = {
val spark = SparkSession
.builder()
.appName(this.getClass.getName)
.master("local[*]")
.getOrCreate()
val sc = spark.sparkContext
val rdd01 = sc.makeRDD(List(1,2,3,4,5), 2)
val rdd02 = sc.makeRDD(List(4,5,6,7,8), 2)
//交集
val rdd03A = rdd01.intersection(rdd02)
println(rdd03A.collect.mkString(",")) //4,5
//并集--不去重
val rdd03B = rdd01.union(rdd02)
println(rdd03B.collect.mkString(",")) //1,2,3,4,5,4,5,6,7,8
//差集
val rdd03C = rdd01.subtract(rdd02)
println(rdd03C.collect.mkString(",")) //2,1,3
//拉链--分区内元素个数相同
//org.apache.spark.SparkException: Can only zip RDDs with same number of elements in each partition
val rdd03D = rdd01.zip(rdd02)
println(rdd03D.collect.mkString(",")) //(1,4),(2,5),(3,6),(4,7),(5,8)
sc.stop()
}
}
11.partitionBy
package sparkCore.rddTransform
import org.apache.spark.HashPartitioner
import org.apache.spark.sql.SparkSession
object RDDTransformV11 {
def main(args: Array[String]): Unit = {
val spark = SparkSession
.builder()
.appName(this.getClass.getName)
.master("local[*]")
.getOrCreate()
val sc = spark.sparkContext
val rdd01 = sc.makeRDD(List(("AA",1),("AA",2),("BB",3),("CC",4)), 3)
val rdd01A = rdd01.partitionBy(new HashPartitioner(rdd01.getNumPartitions))
val rdd01B = rdd01A.mapPartitionsWithIndex {
(index, iter) => {
List((index, iter.toList)).iterator
}
}
println(rdd01B.collect.toBuffer) //ArrayBuffer((0,List((BB,3))), (1,List((AA,1), (AA,2))), (2,List((CC,4))))
sc.stop()
}
}
12.reduceByKey、groupByKey
package sparkCore.rddTransform
import org.apache.spark.rdd.RDD
import org.apache.spark.sql.SparkSession
object RDDTransformV12 {
def main(args: Array[String]): Unit = {
val spark = SparkSession
.builder()
.appName(this.getClass.getName)
.master("local[*]")
.getOrCreate()
val sc = spark.sparkContext
val rdd01 = sc.makeRDD(List(("AA",1),("AA",2),("BB",3),("CC",4)), 3)
//相同key的数据进行value值的聚合操作
val rdd01A = rdd01.reduceByKey((v1,v2) => v1 + v2)
println(rdd01A.collect.toBuffer) //ArrayBuffer((BB,3), (AA,3), (CC,4))
//相同key的数据进行value值的分组操作
val rdd01B: RDD[(String, Iterable[Int])] = rdd01.groupByKey()
//按照一定规则分组后,会将数据整体放到一个组中
val rdd01C: RDD[(String, Iterable[(String, Int)])] = rdd01.groupBy(_._1)
sc.stop()
}
}
13.aggregateByKey、foldByKey
package sparkCore.rddTransform
import org.apache.spark.rdd.RDD
import org.apache.spark.sql.SparkSession
object RDDTransformV13 {
def main(args: Array[String]): Unit = {
val spark = SparkSession
.builder()
.appName(this.getClass.getName)
.master("local[*]")
.getOrCreate()
val sc = spark.sparkContext
val rdd01 = sc.makeRDD(List(("AA", 2), ("AA", 4), ("AA", 8), ("BB", 3), ("BB", 9), ("BB", 27)), 3)
val rdd01A = rdd01.aggregateByKey(0)(
(v1, v2) => math.max(v1, v2),
(v3, v4) => v3 + v4
)
//分区内相同key的value用:拼接,分区间相同key用#拼接
val rdd01B = rdd01.aggregateByKey("@")(
(v1, v2) => s"$v1:$v2",
(v3, v4) => s"$v3#$v4"
)
//计算相同key的平均值,如:AA--14/3 BB=39/3
//初始值(0,0)的类型确定,将决定最终key对应的value值的类型
val rdd01C: RDD[(String, (Int, Int))] = rdd01.aggregateByKey((0, 0))(
(tup, v) => (tup._1 + 1, tup._2 + v),
(tup1, tup2) => (tup1._1 + tup2._1, tup1._2 + tup2._2)
)
val rdd01D = rdd01C.map {
x => {
x match {
case tup: (String, (Int, Int)) => (tup._1, tup._2._2.toDouble / tup._2._1.toDouble)
}
}
}
val rdd01E = rdd01C.map(tup => (tup._1, tup._2._2 / tup._2._1.toDouble))
println(rdd01A.collect.toBuffer) //ArrayBuffer((BB,30), (AA,12))
println(rdd01B.collect.toBuffer) //ArrayBuffer((BB,:3#:9:27), (AA,:2:4#:8))
println(rdd01C.collect.toBuffer) //ArrayBuffer((BB,(3,39)), (AA,(3,14)))
println(rdd01D.collect.toBuffer) //ArrayBuffer((BB,13.0), (AA,4.666666666666667))
println(rdd01E.collect.toBuffer) //ArrayBuffer((BB,13.0), (AA,4.666666666666667))
//分区内和分区间聚合逻辑相同时使用
val rdd02A = rdd01.foldByKey(0)(_ + _)
println(rdd02A.collect.toBuffer) //ArrayBuffer((BB,39), (AA,14))
sc.stop()
}
}
14.combineByKey
package sparkCore.rddTransform
import org.apache.spark.sql.SparkSession
object RDDTransformV14 {
def main(args: Array[String]): Unit = {
val spark = SparkSession
.builder()
.appName(this.getClass.getName)
.master("local[*]")
.getOrCreate()
val sc = spark.sparkContext
val rdd01 = sc.makeRDD(List(("AA", 2), ("AA", 4), ("AA", 8), ("BB", 3), ("BB", 9), ("BB", 27)), 3)
val rdd01A = rdd01.combineByKey(
v1 => (v1, 1),
(tup: (Int, Int), v: Int) => (tup._1 + v, tup._2 + 1),
(tup1: (Int, Int), tup2: (Int, Int)) => (tup1._1 + tup2._1, tup1._2 + tup2._2)
)
val rdd01B = rdd01A.map(tup => (tup._1, tup._2._1 / tup._2._2.toDouble))
println(rdd01A.collect.toBuffer) //ArrayBuffer((BB,(39,3)), (AA,(14,3)))
println(rdd01B.collect.toBuffer) //ArrayBuffer((BB,13.0), (AA,4.666666666666667))
sc.stop()
}
}
15.join、leftOuterJoin
package sparkCore.rddTransform
import org.apache.spark.sql.SparkSession
object RDDTransformV15 {
def main(args: Array[String]): Unit = {
val spark = SparkSession
.builder()
.appName(this.getClass.getName)
.master("local[*]")
.getOrCreate()
val sc = spark.sparkContext
val rdd01 = sc.makeRDD(List(("AA", 1), ("BB", 2), ("CC", 3), ("DD", 4)), 2)
val rdd02 = sc.makeRDD(List(("AA", 1.1), ("BB", 2.2), ("CC", 3.3), ("CC", 4.4)), 2)
val rdd01A = rdd01.join(rdd02)
val rdd01B = rdd01.leftOuterJoin(rdd02)
println(rdd01A.collect.toBuffer) //ArrayBuffer((BB,(2,2.2)), (CC,(3,3.3)), (CC,(3,4.4)), (AA,(1,1.1)))
println(rdd01B.collect.toBuffer) //ArrayBuffer((DD,(4,None)), (BB,(2,Some(2.2))), (CC,(3,Some(3.3))), (CC,(3,Some(4.4))), (AA,(1,Some(1.1))))
sc.stop()
}
}
16.cogroup
package sparkCore.rddTransform
import org.apache.spark.rdd.RDD
import org.apache.spark.sql.SparkSession
object RDDTransformV16 {
def main(args: Array[String]): Unit = {
val spark = SparkSession
.builder()
.appName(this.getClass.getName)
.master("local[*]")
.getOrCreate()
val sc = spark.sparkContext
val rdd01 = sc.makeRDD(List(("AA", 1), ("BB", 2), ("CC", 3), ("DD", 4)), 2)
val rdd02 = sc.makeRDD(List(("AA", 1.1), ("BB", 2.2), ("CC", 3.3), ("CC", 4.4)), 2)
val rdd01A: RDD[(String, (Iterable[Int], Iterable[Double]))] = rdd01.cogroup(rdd02)
rdd01A.collect.foreach(println(_))
// (DD,(CompactBuffer(4),CompactBuffer()))
// (BB,(CompactBuffer(2),CompactBuffer(2.2)))
// (CC,(CompactBuffer(3),CompactBuffer(3.3, 4.4)))
// (AA,(CompactBuffer(1),CompactBuffer(1.1)))
sc.stop()
}
}
二、行动算子
package sparkCore.rddTransform
import org.apache.spark.sql.SparkSession
object RDDActionV1 {
def main(args: Array[String]): Unit = {
val spark = SparkSession
.builder()
.appName(this.getClass.getName)
.master("local[*]")
.getOrCreate()
val sc = spark.sparkContext
val rdd01 = sc.makeRDD(List(1,2,3,4,5), 2)
//行动算子--触发整个job的执行
//按照分区顺序将不同分区的数据采集到Driver端的内存中
val collect: Array[Int] = rdd01.collect
val reduce: Int = rdd01.reduce(_+_)
//数据源中数据的个数
val count: Long = rdd01.count()
//获取数据源中第一个数据
val first: Int = rdd01.first()
//取数据源中的前多少条
val topN: Array[Int] = rdd01.take(3)
//aggregate--初始值会先参与分区内的聚合操作,然后分区间的聚合操作也会再次参与
//aggregateByKey----初始值只会参与分区间的聚合操作
val aggregate = rdd01.aggregate(10)(_+_,_+_) //输出45
//表示分区间和分区内聚合操作逻辑相同
val fold = rdd01.fold(10)(_+_) //输出45
val rdd02 = sc.makeRDD(List(1,2,2,3,3,3), 2)
//统计数据源中各个元素的个数
val map: collection.Map[Int, Long] = rdd02.countByValue()
println(map) //Map(2 -> 2, 1 -> 1, 3 -> 3)
val rdd03 = sc.makeRDD(List(("A", 1), ("A", 2), ("B", 3), ("C", 4),("C",4)),4)
//统计数据源中key的个数
val map2: collection.Map[String, Long] = rdd03.countByKey()
val map3 = rdd03.countByValue()
println(map2) //Map(A -> 2, B -> 1, C -> 2)
println(map3) //Map((A,1) -> 1, (B,3) -> 1, (C,4) -> 2, (A,2) -> 1)
//foreach返回值为()
//RDD中的方法是将计算逻辑发送到Executor端(分布式节点)执行
//Scala集合对象的方法是在同一个节点的内存中完成的
//为了区分不同的处理效果,将RDD的方法称为算子
rdd03.foreach(println(_))
//在项目路径下创建一个output目录,将RDD内的元素按照textfile格式存储进去
//rdd03.repartition(2).saveAsTextFile("output")
sc.stop()
}
}
三、WordCount实现
package sparkCore.rddTransform
import org.apache.spark.rdd.RDD
import org.apache.spark.sql.SparkSession
object WordCountV1 {
def main(args: Array[String]): Unit = {
val spark = SparkSession
.builder()
.appName(this.getClass.getName)
.master("local[*]")
.getOrCreate()
val sc = spark.sparkContext
val lineRDD: RDD[String] = sc.makeRDD(List("hello java", "hello scala", "hello spark"))
//wordcount实现方式一
val wordcount1: RDD[(String, Int)] = lineRDD.flatMap(_.split(" ")).map((_, 1)).groupByKey().map(tuple => (tuple._1, tuple._2.sum))
println("wordcount1: " + wordcount1.collect().toBuffer)
//wordcount实现方式二
val wordcount2: RDD[(String, Int)] = lineRDD.flatMap(_.split(" ")).map((_, 1)).groupByKey().mapValues(it => it.sum)
println("wordcount2: " + wordcount2.collect().toBuffer)
//wordcount实现方式三
val wordcount3: RDD[(String, Int)] = lineRDD.flatMap(_.split(" ")).map((_, 1)).reduceByKey(_ + _)
println("wordcount3: " + wordcount3.collect().toBuffer)
//wordcount实现方式四
val wordcount4: RDD[(String, Int)] = lineRDD.flatMap(_.split(" ")).map((_, 1)).foldByKey(0)(_ + _)
println("wordcount4: " + wordcount4.collect().toBuffer)
//wordcount实现方式五
val wordcount5: RDD[(String, Int)] = lineRDD.flatMap(_.split(" ")).map((_, 1)).combineByKey(x => x, (a: Int, b: Int) => a + b, (m: Int, n: Int) => m + n)
println("wordcount5: " + wordcount5.collect().toBuffer)
//wordcount实现方式六
val wordcount6: RDD[(String, Int)] = lineRDD.flatMap(_.split(" ")).map((_, 1)).aggregateByKey(0)((x: Int, y) => x + y, (m: Int, n: Int) => m + n)
println("wordcount6: " + wordcount6.collect().toBuffer)
//wordcount实现方式七
val wordcount7: collection.Map[String, Long] = lineRDD.flatMap(_.split(" ")).countByValue()
println("wordcount7: " + wordcount7)
//wordcount实现方式八
val wordcount8: collection.Map[String, Long] = lineRDD.flatMap(_.split(" ")).map((_, 1)).countByKey()
println("wordcount8: " + wordcount8)
sc.stop()
}
}