栏目分类:
子分类:
返回
名师互学网用户登录
快速导航关闭
当前搜索
当前分类
子分类
实用工具
热门搜索
名师互学网 > IT > 前沿技术 > 云计算 > 云平台

scala智能化菜品推荐建立推荐模型

云平台 更新时间: 发布时间: IT归档 最新发布 模块sitemap 名妆网 法律咨询 聚返吧 英语巴士网 伯小乐 网商动力

scala智能化菜品推荐建立推荐模型

§能力目标

  1. 能够以基于用户的协同过滤算法建模
  2. 能够对基于物品的协同过滤算法建模
  3. 能够对基于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/testRatings
1.以基于用户的协同过滤算法建模

基于用户的协同过滤,就是通过不同用户对物品的评分来评测用户之间的相似性,搜索目标用户的最近邻用户,然后根据最近邻用户对物品的评分向目标用户进行推荐。其具体实现过程可以描述如下。

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()
  }
}

转载请注明:文章转载自 www.mshxw.com
本文地址:https://www.mshxw.com/it/895585.html
我们一直用心在做
关于我们 文章归档 网站地图 联系我们

版权所有 (c)2021-2022 MSHXW.COM

ICP备案号:晋ICP备2021003244-6号