基于Spark(Scala)淘宝用户行为数据分析
目录
从分布式文件系统HDFS中加载数据、将RDD转换为DataFrame。
统计销售排行前10的商品类目所包含的商品中销售量排行前十的商品
数据集说明
| 名称 | 说明 |
|---|---|
| 用户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)
DAMO开发者矩阵,由阿里巴巴达摩院和中国互联网协会联合发起,致力于探讨最前沿的技术趋势与应用成果,搭建高质量的交流与分享平台,推动技术创新与产业应用链接,围绕“人工智能与新型计算”构建开放共享的开发者生态。
更多推荐


所有评论(0)