目录

数据集说明

代码解析

从分布式文件系统HDFS中加载数据、将RDD转换为DataFrame。

筛选出用户购买数量前十的商品类目

获取排行前10的商品类目所包含的每个商品销售量

统计销售排行前10的商品类目所包含的商品中销售量排行前十的商品

创建一个统计结果的实体类CateAndItemClass

把统计结果的每一行分别存入数组内

编写一个把数组转换为json数据,并保存到本地的方法


 

数据集说明

名称说明
用户ID整数类型,序列化后的用户ID
商品ID整数类型,序列化后的商品ID
商品类目ID整数类型,序列化后的商品所属类目ID
行为类型字符串,枚举类型,包括('pv', 'buy', 'cart', 'fav')
行为发生的时间行为发生的时间戳
行为类型说明
pv商品详情页pv,等价于点击
buy商品购买
cart将商品加入购物车
fav收藏商品

代码解析

从分布式文件系统HDFS中加载数据、将RDD转换为DataFrame。

从RDD转换得到DataFrame有两种方式,一是利用反射机制推断RDD模式,二是使用编程方式定义RDD模式。这里使用第一种方式将RDD转换为DataFrame,需要注意的是应提前定义case class,这样才能被spark隐式转换为DataFrame。

case class Info(userId: Integer, itemId: Integer, cateId: Integer ,action: String, time:String)
    al conf = new SparkConf().setAppName("cate").setMaster("local")
    val sc = new SparkContext(conf)
    val spark = SparkSession.builder().getOrCreate()
    val data1 = "hdfs://192.168.17.200:9000/GraduationData/New_UserBehavior.csv"
    val UserBehavior = sc.textFile(data1).map(_.split(","))
      .map(x => Info(x(0).trim.toInt, x(1).trim.toInt, x(2).trim.toInt, x(3), x(4)))
    val data = spark.createDataFrame(UserBehavior)

筛选出用户购买数量前十的商品类目

首先数据中筛选出"action"字段为"buy"的数据,然后按照"cateId"进行分组并计数,最后按照计数降序排序并取前10条记录,同时重命名"cateId"字段为"cateId1"。

    val top_10_cate = data.filter(data("action") === "buy")
      .groupBy("cateId").count()
      .orderBy(desc("count"))
      .limit(10)
      .withColumnRenamed("cateId", "cateId1")

获取排行前10的商品类目所包含的每个商品销售量

使用join筛选出排行前十的商品类目所包含的商品数据,然后再使用where筛选出用户购买的行为最后计算排行前10的商品类目所包含的每个商品销售量,同时新建一个临时视图“car”。

    
    val ar = top_10_cate.join(data, expr("cateId = cateId1"))
      .where("action like 'buy'")
      .groupBy("cateId", "itemId").count()
ar.createTempView("car")

统计销售排行前10的商品类目所包含的商品中销售量排行前十的商品

使用窗口函数dense_rank()来对数据进行分组,并按某个字段进行排序,然后只取排名在前10的记录。

根据"cateId"字段进行分组,并按照“count“字段的计数进行降序排序,然后取排名在前10的记录。

    val rank = spark.sql("select  *, dense_rank() over (partition by cateId order by count desc) as rank_no from car")
      .where("rank_no <=10").collect()

创建一个统计结果的实体类CateAndItemClass

因为我们要把结果以json数据类型保存到本地,所以我们要先创一个统计结果的实体类以便于我们把结果转化为json类型。

package Dao;


public class CateAndItemClass {
    private String cateId;
    private String itemId;
    private Integer count;
    private Integer rank_no;

    public String getCateId() {
        return cateId;
    }

    public void setCateId(String cateId) {
        this.cateId = cateId;
    }

    public Integer getCount() {
        return count;
    }

    public void setCount(Integer count) {
        this.count = count;
    }

    public Integer getRank_no() {
        return rank_no;
    }

    public void setRank_no(Integer rank_no) {
        this.rank_no = rank_no;
    }

    public String getItemId() {
        return itemId;
    }

    public void setItemId(String itemId) {
        this.itemId = itemId;
    }

    public CateAndItemClass(String cateId, String itemId, Integer count, Integer rank_no) {
        this.cateId = cateId;
        this.itemId = itemId;
        this.count = count;
        this.rank_no = rank_no;
    }
}

把统计结果的每一行分别存入数组内

注意:新建对象的参数类型要与实体类的一致

    val cate_item_list = new Array[CateAndItemClass](956)
    var x = 0
    for (a <- rank) {
      val by_Time = new CateAndItemClass(a(0).toString,a(1).toString,a(2).toString.toInt,a(3).toString.toInt)
      cate_item_list(x) = by_Time
      x = x + 1
    }

编写一个把数组转换为json数据,并保存到本地的方法

因为我们每个指标都有一个实体类,所以这里的writerFile的参数一个泛型参数,以便于引用方法。

  //  把统计数据转换为json格式并保存到本地
  def writerFile[T](a: Array[T], b: String):Unit ={
  val gson = new Gson()
  val json = gson.toJson(a);
  val writer = new BufferedWriter(new FileWriter(b))
  writer.write(json)
  writer.close()
  println(json);
  }

最后引用writerFile把数组转换为json数据类型,并保存到本地

    val file = "../result5.json"
    writerFile(cate_item_list, file)

 

 

 

Logo

DAMO开发者矩阵,由阿里巴巴达摩院和中国互联网协会联合发起,致力于探讨最前沿的技术趋势与应用成果,搭建高质量的交流与分享平台,推动技术创新与产业应用链接,围绕“人工智能与新型计算”构建开放共享的开发者生态。

更多推荐