• Spark案例实际操作


    目录

    数据说明

    需求1:Top10热门品类

    1.1 需求说明

    1.2 实现方案一

    1.2.1 需求分析

    1.2.2 需求实现

    1.2.3 问题

    1.3 实现方案二

    1.3.1 需求分析

    1.3.2 需求实现---结局方案一之中的问题

    1.3.3 问题

    1.4 实现方案三

    1.4.1 需求分析

    1.4.2 需求实现

    1.5 实现方案四

    1.5.1 需求分析

    1.5.2 需求实现

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

    2.1 需求说明

    2.2需求分析

    2.3 功能实现(这里代码之中的架构需要进行注意)

    HotCataegoryTop10的代码 

    HotCategoryTop10SessionController

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

    3.1 需求说明

    3.2 需求分析


    数据说明

    • 上面的黑色底框的数据对应下面的白色底框之中的数据,只是分隔符是不一样的

    说明:上面的数据图是从数据文件中截取的一部分内容,表示为电商网站的用户行为数据,主要包含用户的4种行为:搜索,点击,下单,支付数据规则如下:

    • 数据文件中每行数据采用下划线分隔数据
    • 每一行数据表示用户的一次行为,这个行为只能是4种行为的一种
    • 如果搜索关键字为null,表示数据不是搜索数据
    • 如果点击的品类ID和产品ID为-1,表示数据不是点击数据
    • 针对于下单行为,一次可以下单多个商品,所以品类ID和产品ID可以是多个,id之间采用逗号分隔,如果本次不是下单行为,则数据采用null表示
    • 支付行为和下单行为类似

    详细字段说明:

    样品类:

    1. //用户访问动作表
    2. case class UserVisitAction(
    3. date: String,//用户点击行为的日期
    4. user_id: Long,//用户的ID
    5. session_id: String,//Session的ID
    6. page_id: Long,//某个页面的ID
    7. action_time: String,//动作的时间点
    8. search_keyword: String,//用户搜索的关键词
    9. click_category_id: Long,//某一个商品品类的ID
    10. click_product_id: Long,//某一个商品的ID
    11. order_category_ids: String,//一次订单中所有品类的ID集合
    12. order_product_ids: String,//一次订单中所有商品的ID集合
    13. pay_category_ids: String,//一次支付中所有品类的ID集合
    14. pay_product_ids: String,//一次支付中所有商品的ID集合
    15. city_id: Long
    16. )//城市 id

    需求1:Top10热门品类

    1.1 需求说明

    品类是指产品的分类,大型电商网站品类分多级,咱们的项目中品类只有一级,不同的公司可能对热门的定义不一样。我们按照每个品类的点击、下单、支付的量来统计热门品类。

    鞋 点击数 下单数  支付数

    衣服 点击数 下单数  支付数

    电脑 点击数 下单数  支付数

    例如,综合排名 = 点击数*20%+下单数*30%+支付数*50%

    本项目需求优化为:先按照点击数排名,靠前的就排名高;如果点击数相同,再比较下单数;下单数再相同,就比较支付数。

    1.2 实现方案一

    1.2.1 需求分析

    分别统计每个品类点击的次数,下单的次数和支付的次数:

    (品类,点击总数)(品类,下单总数)(品类,支付总数)

    1.2.2 需求实现

    • 奇怪的面试题,两个杯子,分别是五升水和三升水的杯子,如何得到四升水??

    方式一:5+5-(3+3)=4.首先是将五升水的杯子之中的水倒入三升水空杯子之中,倒出三升水杯子之中的水,将五升水杯子之中剩余的两升水倒入三升水杯子之中。灌满五升水的杯子,将五升水的杯子倒入只有两升水的另一个杯子,剩下的就是四升水。

    方式二:3+3+3-5=4,首先将三升水倒入空的五升水的水杯,在将三升水灌满水倒入五升水的水杯,此刻三升水的水杯剩下一升水,将五升水的水杯之中的水倒掉,三升水的水杯之中的水转移到五生水的水杯之中,再将三升水导入五升水的水杯之中,得到四升水.

    • 补充:如何自动调出RDD的类型,设置步骤是如下所示:
    • 需求实现代码:
      1. package req
      2. import org.apache.spark.rdd.RDD
      3. import org.apache.spark.{SparkConf, SparkContext}
      4. /**
      5. * @anthor Yang
      6. * @create 2022-07-21 22:47
      7. */
      8. object Spark01_Req_HotCategoryTop10 {
      9. def main(args: Array[String]): Unit = {
      10. val conf= new SparkConf().setMaster("local[*]").setAppName("HotCategoryTop10")
      11. val sc = new SparkContext(conf)
      12. //TODO 需求一:Top10热门品类
      13. //TODO 获取原始文件
      14. val fileDatas = sc.textFile("data/user_visit_action.txt")
      15. //TODO 统计分析之前需要将不需要的数据进行过滤-先保留所有的统计数据-对点击数据进行统计
      16. val clickDatas = fileDatas.filter(
      17. data => {
      18. val datas = data.split("_")
      19. val cid = datas(6)
      20. cid != "-1"
      21. }
      22. )
      23. //TODO 统计品类点击数量
      24. val clickCntDatas = clickDatas.map(
      25. data => {
      26. val datas = data.split("_")
      27. val cid = datas(6)
      28. (cid,1)
      29. }
      30. ).reduceByKey(_+_)
      31. //TODO 统计品类下单数量(方法是和上面差不多的)
      32. val orderDatas = fileDatas.filter(
      33. data => {
      34. val datas = data.split("_")
      35. val cid = datas(8)
      36. cid != "null"
      37. }
      38. )
      39. //(1,2,3,4)->(1,1) (2,1) (3.1) 整体变成个体是使用扁平化
      40. //TODO 整体变成部分就是进行一个扁平化的过程
      41. val orderCntDatas = orderDatas.flatMap(
      42. data => {
      43. val datas = data.split("_")
      44. val cid = datas(8)
      45. val cids = cid.split(",")
      46. cids.map((_,1))
      47. }
      48. ).reduceByKey(_+_)
      49. //TODO 统计品类支付数量
      50. val payDatas = fileDatas.filter(
      51. data => {
      52. val datas = data.split("_")
      53. val cid = datas(10)
      54. cid != "null"
      55. }
      56. )
      57. //(1,2,3,4)->(1,1) (2,1) (3.1) 整体变成个体是使用扁平化
      58. //TODO 整体变成部分就是进行一个扁平化的过程
      59. val payCntDatas = payDatas.flatMap(
      60. data => {
      61. val datas = data.split("_")
      62. val cid = datas(10)
      63. val cids = cid.split(",")
      64. cids.map((_,1))
      65. }
      66. ).reduceByKey(_+_)
      67. //TODO 对统计结果进行排序 =》 点击,下单,支付(三个指标全部进行排序完成之后进行选取 -> Tuple)
      68. //val clickSortedDatas = clickCntDatas.sortBy(_._2,false) //这里不是用sortByKey,sortByKet只是对于Key进行一个排序的过程,false表示是降序的过程
      69. //TODO 将点击-》Tuple(品类ID,点击)
      70. //TODO 下单-》Tuple(品类ID,下单)
      71. //TODO 支付-》Tuple(品类ID,支付)
      72. //TODO 将上述的数据变成(品类ID,(点击,下单,支付))
      73. //上述的方式是不能采用JOIN的方式,因为有的数据可能是空的,join要求两边都是有数据的,join的底层是cogroup
      74. //这里不使用下面这一条语句的原因是由于括号的原因
      75. // val value: RDD[(String, (Option[(Option[Int], Option[Int])], Option[Int]))] = clickCntDatas.fullOuterJoin(orderCntDatas).fullOuterJoin(payCntDatas)
      76. val value: RDD[(String, (Iterable[Int], Iterable[Int], Iterable[Int]))] = clickCntDatas.cogroup(orderCntDatas,payCntDatas)
      77. val mapDatas = value.map{
      78. case(cid,(clickIter,orderIter,payIter)) =>
      79. {
      80. var clickcnt =0
      81. var ordercnt =0
      82. var paycnt = 0
      83. val iterator1: Iterator[Int] = clickIter.iterator
      84. if(iterator1.hasNext)
      85. {
      86. clickcnt = iterator1.next()
      87. }
      88. val iterator2: Iterator[Int] = orderIter.iterator
      89. if(iterator2.hasNext)
      90. {
      91. ordercnt = iterator2.next()
      92. }
      93. val iterator3: Iterator[Int] = payIter.iterator
      94. if(iterator3.hasNext)
      95. {
      96. paycnt = iterator3.next()
      97. }
      98. (cid,(clickcnt,ordercnt,paycnt))
      99. }
      100. }
      101. val top10: Array[(String, (Int, Int, Int))] = mapDatas.sortBy(_._2,false).take(10)
      102. //TODO 将统计采集之后的结果打印到控制台上
      103. top10.foreach(println)
      104. sc.stop()
      105. }
      106. }

    1.2.3 问题

    1.三个过滤,同一个RDD(过滤)的重复使用.RDD是不保留数据的,重复使用浪费性能.

    2.cogroup可能存在笛卡尔乘积的问题,相同的时候出现OneToOneDependency,不相同的时候可能会出现ShuffleDependency的现象.------性能瓶颈(可能会出现问题),十几万数据如果要是出现笛卡尔乘积,算子性能会低下.

    1.3 实现方案二

    1.3.1 需求分析

    一次性统计每个品类点击的次数,下单的次数和支付的次数:

    (品类,(点击总数,下单总数,支付总数))

    1.3.2 需求实现---结局方案一之中的问题

    采用reduceByKey的方式代替cogroup,避免出现笛卡尔乘积的问题.

    1. package req
    2. import org.apache.spark.rdd.RDD
    3. import org.apache.spark.{SparkConf, SparkContext}
    4. /**
    5. * @anthor Yang
    6. * @create 2022-07-21 22:47
    7. */
    8. object Spark02_Req_HotCategoryTop10 {
    9. def main(args: Array[String]): Unit = {
    10. val conf= new SparkConf().setMaster("local[*]").setAppName("HotCategoryTop10")
    11. val sc = new SparkContext(conf)
    12. //TODO 需求一:Top10热门品类
    13. //TODO 获取原始文件
    14. val fileDatas = sc.textFile("data/user_visit_action.txt")
    15. //TODO 统计分析之前需要将不需要的数据进行过滤-先保留所有的统计数据-对点击数据进行统计
    16. val clickDatas = fileDatas.filter(
    17. data => {
    18. val datas = data.split("_")
    19. val cid = datas(6)
    20. cid != "-1"
    21. }
    22. )
    23. //TODO 统计品类点击数量
    24. val clickCntDatas = clickDatas.map(
    25. data => {
    26. val datas = data.split("_")
    27. val cid = datas(6)
    28. (cid,1)
    29. }
    30. ).reduceByKey(_+_)
    31. //TODO 统计品类下单数量(方法是和上面差不多的)
    32. val orderDatas = fileDatas.filter(
    33. data => {
    34. val datas = data.split("_")
    35. val cid = datas(8)
    36. cid != "null"
    37. }
    38. )
    39. //(1,2,3,4)->(1,1) (2,1) (3.1) 整体变成个体是使用扁平化
    40. //TODO 整体变成部分就是进行一个扁平化的过程
    41. val orderCntDatas = orderDatas.flatMap(
    42. data => {
    43. val datas = data.split("_")
    44. val cid = datas(8)
    45. val cids = cid.split(",")
    46. cids.map((_,1))
    47. }
    48. ).reduceByKey(_+_)
    49. //TODO 统计品类支付数量
    50. val payDatas = fileDatas.filter(
    51. data => {
    52. val datas = data.split("_")
    53. val cid = datas(10)
    54. cid != "null"
    55. }
    56. )
    57. //(1,2,3,4)->(1,1) (2,1) (3.1) 整体变成个体是使用扁平化
    58. //TODO 整体变成部分就是进行一个扁平化的过程
    59. val payCntDatas = payDatas.flatMap(
    60. data => {
    61. val datas = data.split("_")
    62. val cid = datas(10)
    63. val cids = cid.split(",")
    64. cids.map((_,1))
    65. }
    66. ).reduceByKey(_+_)
    67. //TODO 对统计结果进行排序 =》 点击,下单,支付(三个指标全部进行排序完成之后进行选取 -> Tuple)
    68. //val clickSortedDatas = clickCntDatas.sortBy(_._2,false) //这里不是用sortByKey,sortByKet只是对于Key进行一个排序的过程,false表示是降序的过程
    69. //TODO 将点击-》Tuple(品类ID,点击) ---(品类ID,(点击,0,0))
    70. //TODO 下单-》Tuple(品类ID,下单) ---(品类ID,(0,下单,0))
    71. //TODO 支付-》Tuple(品类ID,支付) ---(品类ID,(0,0,支付))
    72. //TODO 将上述的数据变成(品类ID,(点击,下单,支付))
    73. //上述的方式是不能采用JOIN的方式,因为有的数据可能是空的,join要求两边都是有数据的,join的底层是cogroup
    74. //这里不使用下面这一条语句的原因是由于括号的原因
    75. // val value: RDD[(String, (Option[(Option[Int], Option[Int])], Option[Int]))] = clickCntDatas.fullOuterJoin(orderCntDatas).fullOuterJoin(payCntDatas)
    76. //TODO reduceByKey存在shuffle,但是不存在笛卡尔乘积----修改的地方
    77. val clickMapDatas = clickCntDatas.map{
    78. case( cid,clickCnt ) => {
    79. (cid , (clickCnt,0,0))
    80. }
    81. }
    82. val orderkMapDatas = orderCntDatas.map{
    83. case( cid,orderCnt ) => {
    84. (cid , (0,orderCnt,0))
    85. }
    86. }
    87. val payMapDatas = payCntDatas.map{
    88. case( cid,payCnt ) => {
    89. (cid , (0,0,payCnt))
    90. }
    91. }
    92. //将上述的三个数据进行合并,使用union
    93. val unionRDD: RDD[(String, (Int, Int, Int))] = clickMapDatas.union(orderkMapDatas).union(payMapDatas)
    94. //使用reduceByKey
    95. val reduceRDD = unionRDD.reduceByKey(
    96. (t1,t2) =>
    97. {
    98. (t1._1+t2._1 ,t1._2+t2._2 ,t1._3+t2._3)
    99. }
    100. )
    101. val top10 = reduceRDD.sortBy(_._2,false).take(10)
    102. //TODO 将统计采集之后的结果打印到控制台上
    103. top10.foreach(println)
    104. sc.stop()
    105. }
    106. }

    1.3.3 问题

    reduceByKey存在shuffle的问题,这里我们使用了四个reduceByKey,使用的shuffle次数太多,那么落盘的次数就是很多,想要提高性能,就要想办法减少落盘的过程.

    1.4 实现方案三

    1.4.1 需求分析

    对于上述的问题再次改进

    1.4.2 需求实现

    • 没有使用累加器,但是减少了shuffle.过滤的过程
      1. package req
      2. import org.apache.spark.rdd.RDD
      3. import org.apache.spark.{SparkConf, SparkContext}
      4. import org.codehaus.jackson.annotate.JsonTypeInfo.Id
      5. /**
      6. * @anthor Yang
      7. * @create 2022-07-21 22:47
      8. */
      9. object Spark03_Req_HotCategoryTop10 {
      10. def main(args: Array[String]): Unit = {
      11. val conf= new SparkConf().setMaster("local[*]").setAppName("HotCategoryTop10")
      12. val sc = new SparkContext(conf)
      13. //TODO 需求一:Top10热门品类
      14. //TODO 获取原始文件
      15. val fileDatas = sc.textFile("data/user_visit_action.txt")
      16. val flatDatas = fileDatas.flatMap(//flatMap要求返回的是可以迭代的集合
      17. data => {
      18. var datas = data.split("_")
      19. if( datas(6) != "-1" )
      20. {//点击数据的场合
      21. List((datas(6),(1,0,0)))
      22. }
      23. else if(datas(8) != "null")
      24. {//下单的场合
      25. val id = datas(8)
      26. val ids = id.split(",")
      27. ids.map(//map返回就是迭代的集合
      28. id => {
      29. (id,(0,1,0))
      30. }
      31. )
      32. }
      33. else if(datas(10) != "null")
      34. {
      35. val id= datas(10)
      36. val ids = id.split(",")
      37. ids.map(
      38. id=> {
      39. (id,(0,0,1))
      40. }
      41. )
      42. }
      43. else
      44. {
      45. Nil //全量函数的集合
      46. }
      47. }
      48. )
      49. val top10 = flatDatas.reduceByKey(
      50. (t1 , t2) =>
      51. {
      52. (t1._1 + t2._1 ,t1._2+t2._2 ,t1._3+t2._3)
      53. }
      54. ).sortBy(_._2,false).take(10)
      55. //TODO 将统计采集之后的结果打印到控制台上
      56. top10.foreach(println)
      57. sc.stop()
      58. }
      59. }

      1.5 实现方案四

      1.5.1 需求分析

      使用累加器的过程

      1.5.2 需求实现

    • 实现代码:比较难一些
      1. package com.atguigu.bigdata.spark.req
      2. import org.apache.spark.util.AccumulatorV2
      3. import org.apache.spark.{SparkConf, SparkContext}
      4. import scala.collection.mutable
      5. object Spark01_Req_HotCategoryTop10_3 {
      6. def main(args: Array[String]): Unit = {
      7. val conf = new SparkConf().setMaster("local[*]").setAppName("HotCategoryTop10")
      8. val sc = new SparkContext(conf)
      9. val fileDatas = sc.textFile("data/user_visit_action.txt")
      10. // 创建累加器对象
      11. val acc = new HotCategoryAccumulator()
      12. // 注册累加器
      13. sc.register(acc, "HotCategory")
      14. fileDatas.foreach(
      15. data => {
      16. val datas = data.split("_")
      17. if ( datas(6) != "-1" ) {
      18. // 点击的场合
      19. acc.add( (datas(6), "click") )
      20. } else if ( datas(8) != "null" ) {
      21. // 下单的场合
      22. val id = datas(8)
      23. val ids = id.split(",")
      24. ids.foreach(
      25. id => {
      26. acc.add( (id, "order") )
      27. }
      28. )
      29. } else if ( datas(10) != "null" ) {
      30. // 支付的场合
      31. val id = datas(10)
      32. val ids = id.split(",")
      33. ids.foreach(
      34. id => {
      35. acc.add( (id, "pay") )
      36. }
      37. )
      38. }
      39. }
      40. )
      41. // TODO 获取累加器的结果
      42. val resultMap: mutable.Map[String, HotCategoryCount] = acc.value
      43. val top10 = resultMap.map(_._2).toList.sortWith(
      44. (left, right) => {
      45. if ( left.clickCnt > right.clickCnt ) {
      46. true
      47. } else if ( left.clickCnt == right.clickCnt ) {
      48. if ( left.orderCnt > right.orderCnt ) {
      49. true
      50. } else if ( left.orderCnt == right.orderCnt ) {
      51. left.payCnt > right.payCnt
      52. } else {
      53. false
      54. }
      55. } else {
      56. false
      57. }
      58. }
      59. ).take(10)
      60. top10.foreach(println)
      61. sc.stop()
      62. }
      63. case class HotCategoryCount( cid:String, var clickCnt : Int, var orderCnt : Int, var payCnt : Int )
      64. // TODO 自定义热门点击累加器
      65. // 1. 继承AccumulatorV2
      66. // 2. 定义泛型
      67. // IN : (品类ID,行为类型)
      68. // OUT : Map[品类ID, HotCategoryCount]
      69. // 3. 重写方法 (3 + 3)
      70. class HotCategoryAccumulator extends AccumulatorV2[(String, String), mutable.Map[String, HotCategoryCount]]{
      71. private val map = mutable.Map[String, HotCategoryCount]()
      72. override def isZero: Boolean = {
      73. map.isEmpty
      74. }
      75. override def copy(): AccumulatorV2[(String, String), mutable.Map[String, HotCategoryCount]] = {
      76. new HotCategoryAccumulator()
      77. }
      78. override def reset(): Unit = {
      79. map.clear()
      80. }
      81. override def add(v: (String, String)): Unit = {
      82. val (cid, actionType) = v
      83. val hcc: HotCategoryCount = map.getOrElse(cid, HotCategoryCount(cid, 0, 0, 0))
      84. actionType match {
      85. case "click" => hcc.clickCnt += 1
      86. case "order" => hcc.orderCnt += 1
      87. case "pay" => hcc.payCnt += 1
      88. }
      89. map.update(cid, hcc)
      90. }
      91. override def merge(other: AccumulatorV2[(String, String), mutable.Map[String, HotCategoryCount]]): Unit = {
      92. other.value.foreach {
      93. case ( cid, otherHCC ) => {
      94. val thisHCC: HotCategoryCount = map.getOrElse(cid, HotCategoryCount(cid, 0, 0, 0))
      95. thisHCC.clickCnt += otherHCC.clickCnt
      96. thisHCC.orderCnt += otherHCC.orderCnt
      97. thisHCC.payCnt += otherHCC.payCnt
      98. map.update(cid, thisHCC)
      99. }
      100. }
      101. }
      102. override def value: mutable.Map[String, HotCategoryCount] = {
      103. map
      104. }
      105. }
      106. }

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

    2.1 需求说明

    在需求一的基础上,增加每个品类用户session点击统计

    2.2需求分析

    JavaEE        Web:浏览器(zhangsan) => 服务器(lisi) [lisi懵了]       Session:会话(通信状态)

    2.3 功能实现(这里代码之中的架构需要进行注意)

    HotCataegoryTop10的代码 

    • Common之中的代码,用来呈递上面架构之中的架构部分,对里面的代码进行相应的封装:

    ①对于Application部分的封装

    1. package com.atguigu.bigdata.spark.summer.common
    2. import com.atguigu.bigdata.spark.summer.util.EnvCache
    3. import org.apache.spark.{SparkConf, SparkContext}
    4. trait TApplication {
    5. //下面第一个括号之中代表的是参数,第二个代表的是进行执行的操作
    6. def execute(master:String = "local[*]", appName:String)( op: =>Unit ): Unit = {
    7. val conf: SparkConf = new SparkConf().setMaster(master).setAppName(appName)
    8. val sc = new SparkContext(conf)
    9. EnvCache.put(sc) //用来存放sc这个值,后面可以拿来用
    10. try {
    11. op
    12. } catch {
    13. case e: Exception => e.printStackTrace()
    14. }
    15. sc.stop()
    16. EnvCache.clear()
    17. }
    18. }

    ②对于Controller部分的封装

    1. package com.atguigu.bigdata.spark.summer.common
    2. trait TController {
    3. def dispatch(): Unit
    4. }

    ③对于Dao部分的封装

    1. package com.atguigu.bigdata.spark.summer.common
    2. import com.atguigu.bigdata.spark.summer.util.EnvCache
    3. import org.apache.spark.SparkContext
    4. import scala.io.{BufferedSource, Source}
    5. trait TDao {
    6. def readFile( path : String ) = {
    7. // e:/data/word.txt
    8. val source: BufferedSource = Source.fromFile(EnvCache.get() + path)
    9. val lines = source.getLines().toList
    10. source.close()
    11. lines
    12. }
    13. def readFileBySpark( path : String ) = {
    14. //对于前面使用的sc进行一个读取的过程
    15. EnvCache.get().asInstanceOf[SparkContext].textFile(path)
    16. }
    17. }

    ④Service部分的封装

    1. package com.atguigu.bigdata.spark.summer.common
    2. trait TService {
    3. def analysis() : Any = {
    4. }
    5. def analysis( data : Any ) : Any = {
    6. }
    7. }
    • Application(应用的起点)部分的代码:
    1. package com.atguigu.bigdata.spark.summer.application
    2. import com.atguigu.bigdata.spark.summer.common.TApplication
    3. import com.atguigu.bigdata.spark.summer.controller.{HotCategoryTop10Controller, WordCountController}
    4. //这里进行继承App便是可以之间执行的主函数,可以不用写出main函数
    5. object HotCategoryTop10Application extends TApplication with App{
    6. //这里的execute本来是两个参数,但是其中一个是已经定义好的默认参数
    7. execute(appName = "HotCategoryTop10"){
    8. //执行controller
    9. val controller = new HotCategoryTop10Controller
    10. controller.dispatch()
    11. }
    12. }
    • Controller(进行调度)部分的代码:
    1. package com.atguigu.bigdata.spark.summer.controller
    2. import com.atguigu.bigdata.spark.summer.common.TController
    3. import com.atguigu.bigdata.spark.summer.service.HotCategoryTop10Service
    4. class HotCategoryTop10Controller extends TController {
    5. private val hotCategoryTop10Service = new HotCategoryTop10Service
    6. override def dispatch(): Unit = {
    7. val result: Array[(String, (Int, Int, Int))] = hotCategoryTop10Service.analysis()
    8. result.foreach(println)
    9. }
    10. }
    • Dao(跟数据打交道)部分的代码
    1. package com.atguigu.bigdata.spark.summer.dao
    2. import com.atguigu.bigdata.spark.summer.common.TDao
    3. class HotCategoryTop10Dao extends TDao {
    4. }
    • Service(执行逻辑)部分的代码
    1. package com.atguigu.bigdata.spark.summer.service
    2. import com.atguigu.bigdata.spark.summer.common.TService
    3. import com.atguigu.bigdata.spark.summer.dao.HotCategoryTop10Dao
    4. class HotCategoryTop10Service extends TService {
    5. private val hotCategoryTop10Dao = new HotCategoryTop10Dao
    6. override def analysis() = {
    7. val fileDatas = hotCategoryTop10Dao.readFileBySpark("data/user_visit_action.txt")
    8. val flatDatas = fileDatas.flatMap(
    9. data => {
    10. var datas = data.split("_")
    11. if ( datas(6) != "-1" ) {
    12. // 点击数据的场合
    13. List((datas(6), (1, 0, 0)))
    14. } else if ( datas(8) != "null" ) {
    15. // 下单数据的场合
    16. val id = datas(8)
    17. val ids = id.split(",")
    18. ids.map(
    19. id => {
    20. (id, (0, 1, 0))
    21. }
    22. )
    23. } else if ( datas(10) != "null" ) {
    24. // 支付数据的场合
    25. val id = datas(10)
    26. val ids = id.split(",")
    27. ids.map(
    28. id => {
    29. (id, (0, 0, 1))
    30. }
    31. )
    32. } else {
    33. Nil
    34. }
    35. }
    36. )
    37. val top10 = flatDatas.reduceByKey(
    38. (t1, t2) => {
    39. ( t1._1 + t2._1, t1._2 + t2._2, t1._3 + t2._3 )
    40. }
    41. ).sortBy(_._2, false).take(10)
    42. top10
    43. }
    44. }
    • Util(主线程之中用来存放rootPath连接)的代码:
    1. package com.atguigu.bigdata.spark.summer.util
    2. object EnvCache {
    3. private val envCache : ThreadLocal[Object] = new ThreadLocal[Object]
    4. def put( data : Object ): Unit = {
    5. envCache.set(data)
    6. }
    7. def get() = {
    8. envCache.get()
    9. }
    10. def clear(): Unit = {
    11. envCache.remove()
    12. }
    13. }

    HotCategoryTop10SessionController

    这部分的代码是在上面的原有基础上进行一个增加的过程.

    •  Application(应用的起点)部分的代码:
    1. package com.atguigu.bigdata.spark.summer.application
    2. import com.atguigu.bigdata.spark.summer.common.TApplication
    3. import com.atguigu.bigdata.spark.summer.controller.{HotCategoryTop10Controller, HotCategoryTop10SessionController}
    4. object HotCategoryTop10SessionApplication extends TApplication with App{
    5. execute(appName = "HotCategoryTop10Session"){
    6. val controller = new HotCategoryTop10SessionController
    7. controller.dispatch()
    8. }
    9. }
    • Controller(进行调度)部分的代码:
    1. package com.atguigu.bigdata.spark.summer.controller
    2. import com.atguigu.bigdata.spark.summer.common.TController
    3. import com.atguigu.bigdata.spark.summer.service.{HotCategoryTop10Service, HotCategoryTop10SessionService}
    4. class HotCategoryTop10SessionController extends TController {
    5. //由于这里是需要用到上面的逻辑,因此这里是存在两个代码
    6. private val hotCategoryTop10Service = new HotCategoryTop10Service
    7. private val hotCategoryTop10SessionService = new HotCategoryTop10SessionService
    8. override def dispatch(): Unit = {
    9. val top10: Array[(String, (Int, Int, Int))] = hotCategoryTop10Service.analysis()
    10. val result = hotCategoryTop10SessionService.analysis(top10.map(_._1))//只是存取Id
    11. result.foreach(println)
    12. }
    13. }
    • Dao(跟数据打交道)部分的代码
    1. package com.atguigu.bigdata.spark.summer.dao
    2. import com.atguigu.bigdata.spark.summer.common.TDao
    3. class HotCategoryTop10SessionDao extends TDao {
    4. }
    • Service(执行逻辑)部分的代码
    1. package com.atguigu.bigdata.spark.summer.service
    2. import com.atguigu.bigdata.spark.summer.bean.UserVisitAction
    3. import com.atguigu.bigdata.spark.summer.common.TService
    4. import com.atguigu.bigdata.spark.summer.dao.{HotCategoryTop10Dao, HotCategoryTop10SessionDao}
    5. import org.apache.spark.rdd.RDD
    6. class HotCategoryTop10SessionService extends TService {
    7. private val hotCategoryTop10SessionDao = new HotCategoryTop10SessionDao
    8. override def analysis( data : Any ) = {
    9. val topIds: Array[String] = data.asInstanceOf[Array[String]]
    10. val fileDatas = hotCategoryTop10SessionDao.readFileBySpark("data/user_visit_action.txt")
    11. val actionDatas = fileDatas.map(
    12. data => {
    13. val datas = data.split("_")
    14. //将字符串转换为long类型
    15. UserVisitAction(
    16. datas(0),
    17. datas(1).toLong,
    18. datas(2),
    19. datas(3).toLong,
    20. datas(4),
    21. datas(5),
    22. datas(6).toLong,
    23. datas(7).toLong,
    24. datas(8),
    25. datas(9),
    26. datas(10),
    27. datas(11),
    28. datas(12).toLong
    29. )
    30. }
    31. )
    32. val clickDatas = actionDatas.filter {
    33. data => {
    34. if ( data.click_category_id != -1 ) //点击不是-1
    35. {
    36. topIds.contains(data.click_category_id.toString)
    37. } else {
    38. false
    39. }
    40. }
    41. }
    42. val reduceDatas = clickDatas.map(
    43. data => {
    44. (( data.click_category_id, data.session_id ), 1)
    45. }
    46. ).reduceByKey(_+_)
    47. val groupDatas: RDD[(Long, Iterable[(String, Int)])] = reduceDatas.map {
    48. case ((cid, sid), cnt) => {
    49. (cid, (sid, cnt))
    50. }
    51. }.groupByKey()
    52. groupDatas.mapValues(
    53. iter => {
    54. iter.toList.sortBy(_._2)(Ordering.Int.reverse).take(10)
    55. }
    56. ).collect()
    57. }
    58. }
    • Bean之中增加UserVisitAcition代码如下所示:
    1. package com.atguigu.bigdata.spark.summer.bean
    2. //用户访问动作表
    3. case class UserVisitAction(
    4. date: String,//用户点击行为的日期
    5. user_id: Long,//用户的ID
    6. session_id: String,//Session的ID
    7. page_id: Long,//某个页面的ID
    8. action_time: String,//动作的时间点
    9. search_keyword: String,//用户搜索的关键词
    10. click_category_id: Long,//某一个商品品类的ID
    11. click_product_id: Long,//某一个商品的ID
    12. order_category_ids: String,//一次订单中所有品类的ID集合
    13. order_product_ids: String,//一次订单中所有商品的ID集合
    14. pay_category_ids: String,//一次支付中所有品类的ID集合
    15. pay_product_ids: String,//一次支付中所有商品的ID集合
    16. city_id: Long//城市 id
    17. )

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

    3.1 需求说明

    1)页面单跳转化率

    计算页面单跳转化率,什么是页面单跳转换率,比如一个用户在一次 Session 过程中访问的页面路径 3,5,7,9,10,21,那么页面 3 跳到页面 5 叫一次单跳,7-9 也叫一次单跳,那么单跳转化率就是要统计页面点击的概率。

    比如:计算 3-5 的单跳转化率,先获取符合条件的 Session 对于页面 3 的访问次数(PV)为 A,然后获取符合条件的 Session 中访问了页面 3 又紧接着访问了页面 5 的次数为 B,那么 B/A 就是 3-5 的页面单跳转化率。

    为什么IE被淘汰了?因为IE不支持新的技术,而且不能够装入插件.不要说什么产品经理没用,你是可以替换的,但是产品经理是不可以替换的.

    2)统计页面单跳转化率意义

    产品经理和运营总监,可以根据这个指标,去尝试分析,整个网站,产品,各个页面的表现怎么样,是不是需要去优化产品的布局;吸引用户最终可以进入最后的支付页面。

    数据分析师,可以此数据做更深一步的计算和分析。

    企业管理层,可以看到整个公司的网站,各个页面的之间的跳转的表现如何,可以适当调整公司的经营战略或策略。

    信息产生过剩:产生的大量的垃圾信息,在垃圾信息之中,找出有用的信息.

    3.2 需求分析

     第一列是进行一个采集的过程,数据产生乱序的原因,如下流程图所示:

     上述过程是需要进行计算一个分子与分母的过程,分子(首页 - 详情)  分母(首页)

    1. package com.atguigu.bigdata.spark.req
    2. import com.atguigu.bigdata.spark.summer.bean.UserVisitAction
    3. object Spark03_Req_PageFlow {
    4. def main(args: Array[String]): Unit = {
    5. val conf = new SparkConf().setMaster("local[*]").setAppName("Pageflow")
    6. val sc = new SparkContext(conf)
    7. val fileDatas = sc.textFile("data/user_visit_action.txt")
    8. //包装成为样例类
    9. val actionDatas = fileDatas.map(
    10. data => {
    11. val datas = data.split("_")
    12. UserVisitAction(
    13. datas(0),
    14. datas(1).toLong,
    15. datas(2),
    16. datas(3).toLong,
    17. datas(4),
    18. datas(5),
    19. datas(6).toLong,
    20. datas(7).toLong,
    21. datas(8),
    22. datas(9),
    23. datas(10),
    24. datas(11),
    25. datas(12).toLong
    26. )
    27. }
    28. )
    29. actionDatas.cache()//做一个缓存过程
    30. // 【1,2,3,4,5,6,7】
    31. val okIds = List(1,2,3,4,5,6,7)
    32. // 【(1,2),(2,3)】
    33. val okFlowIds = okIds.zip(okIds.tail)
    34. // TODO 分母的计算 Long是页面id,int是统计结果
    35. val result: Map[Long, Int] = actionDatas.filter(
    36. action => {
    37. okIds.init.contains(action.page_id.toInt)
    38. }
    39. ).map(
    40. action => {
    41. (action.page_id, 1)
    42. }
    43. ).reduceByKey(_ + _).collect().toMap
    44. // TODO 分子的计算
    45. // 1. 将数据按照session进行分组
    46. val groupRDD: RDD[(String, Iterable[UserVisitAction])] = actionDatas.groupBy(_.session_id)
    47. // 将分组后的数据进行组内排序
    48. val mapRDD = groupRDD.mapValues(
    49. iter => {
    50. //toList含义就是将上面的东西进行转换,进行可以迭代操作
    51. val actions: List[UserVisitAction] = iter.toList.sortBy(_.action_time)
    52. //【1,2,3,4,5,6,7】
    53. //【2,3,4,5,6,7】
    54. // 滑窗
    55. //【1-2,2-3,3-4,4-5,5-6,6-7】
    56. val ids: List[Int] = actions.map(_.page_id.toInt)
    57. // val iterator: Iterator[List[Long]] = ids.sliding(2)
    58. // while ( iterator.hasNext ) {
    59. // val longs: List[Long] = iterator.next()
    60. // (longs.head, longs.last)
    61. // }
    62. val flowIds: List[(Int, Int)] = ids.zip(ids.tail)
    63. flowIds.filter(
    64. ids => {
    65. okFlowIds.contains(ids)
    66. }
    67. )
    68. }
    69. )
    70. val mapRDD2 = mapRDD.map(_._2)
    71. val flatRDD = mapRDD2.flatMap(list => list)
    72. // 分子计算完毕
    73. val reduceRDD = flatRDD.map((_, 1)).reduceByKey(_ + _)
    74. // TODO 单挑转换率的统计
    75. reduceRDD.foreach {
    76. case ( (id1, id2), cnt ) => {
    77. println(s"页面【${id1}-${id2}】单挑转换率为 :" + ( cnt.toDouble / result.getOrElse(id1, 1) ))
    78. }
    79. }
    80. sc.stop()
    81. }
    82. }
  • 相关阅读:
    Docker - 容器的网络模式
    MySQL 数据库 查询定义参数【模糊查询】
    终于还是熬不住了,转行了,分享一波刚学到的知识吧,字符串的自带函数.py
    (免费领源码)java#SSM#Mysql学院教室管理系统81671-计算机毕业设计项目选题推荐
    vue -权限管理-指令-v-permission
    【MySQL入门实战5】-Linux PRM 包安装MySQL
    Java#14(StringJoiner)
    Vue 全局状态管理工具 pinia
    《Linux运维总结:内网服务器通过代理访问外网服务器(方法一)》
    python迭代器
  • 原文地址:https://blog.csdn.net/m0_47489229/article/details/125922298