文章目录§能力目标
- 能够以基于用户的协同过滤算法建模
- 能够对基于物品的协同过滤算法建模
- 能够对基于Spark ALS的协同过滤算法建模(拓展)
- 前言
- 一、推荐算法的选择
- 1.以基于用户的协同过滤算法建模
- 2.计算相似度
- 3.寻找与目标用户最近邻的K个用户
- 4.通过这K个用户进行推荐
- 二、以基于物品的协同过滤算法建模
- 1.打包执行以上两个创建推荐模型的程序
- 1.1第一个程序
- 1.2 查看结果
- 1.3 第二个程序
- 1.4 查看结果
- 三、以基于Spark ALS算法建模
- 1.ALS推荐算法
- 2.Spark ALS算法
- 3.以Spark ALS算法建模
- 4.建模时的参数寻优
前言
提示:以下是本篇文章正文内容,下面案例可供参考
做此篇文章的项目请自行写完上一篇我的博客文章项目(链接放下面了!!)
Scala智能化菜品推荐数据预处理(点击这里)
一、推荐算法的选择协同过滤算法,包括多种可实现的算法。
本例将选用3种不同的协同过滤算法,分别来建立推荐模型,最后进行综合评估。协同过滤算法包括:
- (1) 基于用户的协同过滤算法。
- (2) 基于物品的协同过滤算法。
- (3) 基于Spark ALS的协同过滤算法。
训练集、验证集、测试集保存路径 trainDataPath: /data/MealRatings/trainRatings validateDataPath:/data/MealRatings/validateRatings testDataPath:/data/MealRatings/testRatings1.以基于用户的协同过滤算法建模
基于用户的协同过滤,就是通过不同用户对物品的评分来评测用户之间的相似性,搜索目标用户的最近邻用户,然后根据最近邻用户对物品的评分向目标用户进行推荐。其具体实现过程可以描述如下。
2.计算相似度用户之间的相似度通过每个用户对物品的评分向量计算得到。相似度的计算可以使用任何向量相似度计算公式,常用的相似度计算公式有Jaccard公式、余弦相似度公式、欧式距离公式等
3.寻找与目标用户最近邻的K个用户在计算出各个用户之间的相似度后,可以找到所有与目标用户的相似度大于某一阈值的近邻用户(初步的粗略过滤),然后对这些用户按照相似度值进行排序,得到前K个近邻用户。
4.通过这K个用户进行推荐在获得K个近邻用户后,怎么推荐呢?这里的方式有多种。比如使用相似度和所有K个用户的物品对应加权进行推荐。
参考上述算法原理,以Spark编程来逐步实现。(1)首先加载训练数据集,为减小用户评分过于稀疏的可能影响,在此可以使用“单用户评价过的最小物品数”对其进行过滤。(2)然后根据用户对物品的评分向量获得用户相似度.(3)最后匹配训练集数据生成推荐模型。具体过程的实现代码如下所示:
package cn.mealdata
import org.apache.spark.{SparkConf, SparkContext}
import scala.math._
object UserBaseModelCreate {
def main(args: Array[String]): Unit = {
if (args.length != 5) {
System.err.println("Usage: cn.mealdata.UserBasedModelCreate "+
" ")
}
val trainDataPath = args(0)
val modelPath=args(1)
val minItemsRatedPerUser=args(2).toInt
val recommendItemNum=args(3).toInt
val splitter=args(4)
val appName=" UserBased CF Create Model "
val conf=new SparkConf().setAppName("appName")
val sc=new SparkContext()
val trainDataRaw=sc.textFile(trainDataPath).map{
x=>val fields=x.slice(1,x.size-1).split(splitter);
(fields(0).toInt,fields(1).toInt,fields(2).toDouble)}
val trainDataFilered=trainDataRaw.groupBy(_._1).
filter(data=>data._2.toList.size >= minItemsRatedPerUser).flatMap(_._2)
val trainUserItemRating=trainDataFilered.map{ case (user,item,rating)=>(user,(item,rating))}
val trainUserRating=trainDataFilered.map{case (user,item,rating)=>(user,rating)}.groupByKey().map{
x=>(x._1,x._2.reduce(_ + _)/x._2.count(x=>true))}
val userItemBase=trainUserItemRating.join(trainUserRating).map(x=>(x._1,x._2._1._1,x._2._1._2,x._2._2))
val itemUserBase=userItemBase.map(x=>(x._2,(x._1,x._3,x._4)))
val itemMatrix=itemUserBase.join(itemUserBase).filter((f=>f._2._1._1 < f._2._2._1))
println("itemMatrix records count : "+ itemMatrix.count)
val userSimilarityBase = itemMatrix.map(
f=>((f._2._1._1,f._2._2._1),(f._2._1._2,f._2._1._3,f._2._2._2,f._2._2._3)))
val userSimilarityPre = userSimilarityBase.map(data=>{
val user1=data._1._1
val user2=data._1._2
val similarity = (min(data._2._1,data._2._3))/(data._2._2+data._2._4)
((user1,user2),similarity)
}).combineByKey(
x=>x,
(x:Double,y:Double)=>(x+y),
(x:Double,y:Double)=>(x+y)).cache()
val userSimilarity1=userSimilarityPre.map(x=>(x._1._1,(x._1._2,x._2)))
val userSimilarity2=userSimilarityPre.map(x=>(x._1._2,(x._1._1,x._2)))
val statisticsPre1=trainUserItemRating.map(x=>(x._1,x._2._1)).join(userSimilarity1).
map(x=>(x._2._2._1,(x._2._1,x._2._2._2))).cache()
val statisticsPre2=trainUserItemRating.map(x=>(x._1,x._2._1)).join(userSimilarity2).
map(x=>(x._2._2._1,(x._2._1,x._2._2._2))).cache()
val statistics=statisticsPre1.union(statisticsPre2).combineByKey(
(x:(Int,Double)) =>List(x),
(c:List[(Int,Double)],x:(Int,Double))=>c:+x,
(c1:List[(Int,Double)],c2:List[(Int,Double)])=>c1:::c2).cache()
val dataModel=statistics.
map(data=>{
val key=data._1;
val value=data._2.sortWith(_._2>_._2);
if (value.size>recommendItemNum){
(key,value.slice(0,recommendItemNum))
} else{
(key,value)
}
}).map(x=>(x._1,x._2.map(x=>x._1)))
println("Model records count : " + dataModel.count)
dataModel.repartition(6).saveAsObjectFile(modelPath)
println("Model saved")
sc.stop()
}
}
将用户相似度模型与训练集数据进行匹配,获得每个用户的K个可推荐物品列表,也可以视为推荐结果集。将推荐结果集存储在HDFS上,后续将进行模型评价
二、以基于物品的协同过滤算法建模基于物品的协同过滤推荐算法,其基本思想是用户对物品的预测评分可以由该用户对与该物品相似度最高的K个邻居物品的评分通过加权平均计算得到,如下图所示,对物品1感兴趣的用户也都对物品2n感兴趣,因此物品1和物品2n的相似度较高,它们属于相似物品,而用户t目前对物品2-n感兴趣,但还没发现物品1,因此可将物品1推荐给用户t。
参考上述算法原理,以Spark 编程来逐步实现。
-
(1)首先加载训练数据集,为减小用户评分过于稀疏的可能影响,在此可以使用“单用户评价过的最小物品数”对其进行过滤。
-
(2)然后根据用户对物品的评分向量获得物品相似度。
-
(3)最后匹配训练集数据生成推荐模型。
具体实现过程如下所示。
package cn.mealdata
import org.apache.spark.{SparkContext,SparkConf}
import scala.math._
object ItemBasedModelCreate {
def main(args: Array[String]): Unit = {
if(args.length !=5){
System.err.println("Usage:com.tipdm.ItemBasedModelCreate " +
" ")
}
val trainDataPath = args(0)
val minRatedNumPerUser=args(2).toInt
val recommendIteNum=args(2).toInt
val modelPath =args(3)
val splitter=args(4)
val appName=" ItemBased CF Create Model"
val conf=new SparkConf().setAppName("appName")
val sc=new SparkContext()
val trainData=sc.textFile(trainDataPath).map{x=> val fields=x.slice(1,x.size-1).split(splitter);
(fields(0).toInt,fields(1).toInt,fields(2).toDouble)}
val trainDataFiltered =trainData.groupBy(_._1).
filter(data=>data._2.toList.size >= minRatedNumPerUser).flatMap(_._2).cache()
val trainUserItemRating =trainData.map{ case (user,item,rating) => (item ,(user,rating))}
val trainItemRating =trainData.map{case (user,item,rating) => (item,rating)}.groupByKey().map{
x=> (x._1,x._2.reduce(_+_)/x._2.count(x=>true))
}
val itemUserBase=trainUserItemRating.join(trainItemRating).
map(x=>(x._2._1._1,(x._1,x._2._1._2,x._2._2))).cache()
val itemMatrix=itemUserBase.join(itemUserBase).filter((f => f._2._1._1 < f._2._2._1))
val itemSimilarityBase = itemMatrix.
map(f => ((f._2._1._1,f._2._2._1),(f._2._1._2,f._2._1._3,f._2._2._2,f._2._2._3)))
val itemSimilarityPre =itemSimilarityBase.map(data => {
val item1=data._1._1
val item2=data._1._2
val similarity =(min(data._2._1,data._2._3))*1.0/(data._2._2+data._2._4)
((item1,item2),similarity)
}).combineByKey(
x=>x,
(x:Double,y:Double) => (x+y),
(x:Double,y:Double) => (x+y))
val itemSimilarity=itemSimilarityPre.map(x=> ((x._1._2,x._1._1),x._2)).union(itemSimilarityPre).
map(x =>(x._1._1,(x._1._2,x._2)))
println("itemSimilarity records count: "+itemSimilarity.count)
val dataModelPre = itemSimilarity.combineByKey(
(x:(Int,Double)) => List(x),
(c:List[(Int,Double)],x:(Int,Double)) => c:+x,
(c1:List[(Int,Double)],c2:List[(Int,Double)]) => c1:::c2)
val dataModel =trainDataFiltered.map(x=>(x._2,x._1)).join(dataModelPre)
val recommendModel =dataModel.flatMap(joined => {
joined._2._2.map(f => (joined._2._1,f._1,f._2))}).sortBy( x=>
(x._1,x._3),ascending = false).
map(x=> (x._1,x._2)).
combineByKey(
(x:Int) => List(x),
(c:List[Int],x:Int) => c:+x,
(c1:List[Int],c2:List[Int])=> c1:::c2).map(x =>
(x._1,x._2.take(recommendIteNum)))
println("Recommend Model count: " + recommendModel.count)
recommendModel.repartition(numPartitions=6).saveAsObjectFile(modelPath)
sc.stop()
}
}
通过计算获得物品之间的相似度模型,与训练集数据进行匹配,获得每个用户的可推荐物品列表,将推荐结果集存储到HDFS上
1.打包执行以上两个创建推荐模型的程序 1.1第一个程序[root@hadoop mealRating]# /opt/spark/bin/spark-submit --master local --class cn.mealdata.UserBasedModelCreate /root/mealRating/HadoopSpark3-1.0-SNAPSHOT.jar /data/MealRatings/trainRatings
object ALSModelOptimize {
val appName = "Optimize ALS Model Parameters"
val conf = new SparkConf().setAppName(appName)
val sc = new SparkContext(conf)
sc.setLogLevel("WARN")
def main(args: Array[String]) = {
if (args.length != 7) {
System.err.println("Usage:ALSModelOptimize requires: 7 input fields ")
}
// 匹配输入参数
val trainDataPath = args(0)
val validateDataPath = args(1)
val listRank = args(2).split(",").map(_.toInt)
val listIteration = args(3).split(",").map(_.toInt)
val listLambda = args(4).split(",").map(_.toDouble)
val paraOutputPath = args(5)
val splitter = args(6)
// 定义计算均方根误差的函数:computeRMSE
def computeRMSE(model:MatrixFactorizationModel, data:RDD[Rating]): Double = {
val usersProducts = data.map(x=>(x.user,x.product))
val ratingsAndPredictions = data.map{case Rating(user,product,rating)=>((user,product),rating)
}.join(model.predict(usersProducts).map{case Rating(user,product,rating)=>((user,product),rating)}).values ;
math.sqrt(ratingsAndPredictions.map(x=>(x._1-x._2)*(x._1-x._2)).mean())}
// 加载训练集数据
val trainData = sc.textFile(trainDataPath).map{x=>val fields=x.slice(1,x.size-1).split(splitter);
(fields(0).toInt,fields(1).toInt,fields(2).toDouble)}
val trainDataRating= trainData.map(x=>Rating(x._1,x._2,x._3))
// 加载验证集数据
val validateData = sc.textFile(validateDataPath).map{x=>val fields=x.slice(1,x.size-1).split(splitter);
(fields(0).toInt,fields(1).toInt,fields(2).toDouble)}
val validateDataRating= validateData.map(x=>Rating(x._1,x._2,x._3))
// 初始化最优参数,取极端值
var bestRMSE = Double.MaxValue
var bestRank = -10
var bestIteration = -10
var bestLambda = -1.0
// 参数寻优
for( rank<- listRank;lambda<-listLambda;iter<-listIteration) {
val model = ALS.train(trainDataRating,rank,iter,lambda);
val validationRMSE = computeRMSE(model,validateDataRating);
if(validationRMSE
bestRMSE=validationRMSE;
bestRank=rank;
bestLambda=lambda;
bestIteration=iter}
}
// 输出最优参数组
println("BestRank:Iteration:BestLambda => BestRMSE")
println(bestRank + ": " + bestIteration + ": " + bestLambda + " => " + bestRMSE)
val result = Array(bestRank + "," + bestIteration + "," + bestLambda)
sc.parallelize(result).repartition(1).saveAsTextFile(paraOutputPath)
sc.stop()
}
}
在获得最优的参数组之后,使用这些参数来建模,执行Sspark ALS中的train方法,具体“以最优参数组建立模型”代码如下:
package cn.mealdata
import org.apache.spark.{SparkContext, SparkConf}
import org.apache.spark.mllib.recommendation.ALS
import org.apache.spark.mllib.recommendation.Rating
object ALSModelCreate {
val appName = "Create ALS Model "
val conf = new SparkConf().setAppName(appName)
val sc = new SparkContext(conf)
sc.setLogLevel("WARN")
def main(args: Array[String]) = {
if (args.length != 6) {
System.err.println("Usage: ALSModelCreate requires: 6 input fields " +
" ")
}
// 匹配输入参数
val trainDataPath = args(0)
val modelPath = args(1)
val rank = args(2).toInt
val iteration = args(3).toInt
val lambda = args(4).toDouble
val splitter = args(5)
// 加载训练集数据
val trainData = sc.textFile(trainDataPath).map{x=>val fields=x.slice(1,x.size-1).split(splitter);
(fields(0).toInt,fields(1).toInt,fields(2).toDouble)}
val trainDataRating= trainData.map(x=>Rating(x._1,x._2,x._3))
// 建立ALS模型
val model = ALS.train(trainDataRating, rank, iteration, lambda)
// 存储ALS模型
model.save(sc,modelPath)
println("Model saved")
sc.stop()
}
}



