spark基于物品的推荐_大数据项目之电影推荐系统(下)

一、实时推荐服务建设
1、实时推荐服务
实时计算与离线计算应用于推荐系统上最大的不同在于实时计算推荐结果应该
反映最近一段时间用户近期的偏好,而离线计算推荐结果则是根据用户从第一次评
分起的所有评分记录来计算用户总体的偏好。
用户对物品的偏好随着时间的推移总是会改变的。比如一个用户 u 在某时刻对
电影 p 给予了极高的评分,那么在近期一段时候,u 极有可能很喜欢与电影 p 类
似的其他电影;而如果用户 u 在某时刻对电影 q 给予了极低的评分,那么在近期一
段时候,u 极有可能不喜欢与电影 q 类似的其他电影。所以对于实时推荐,当用户
对一个电影进行了评价后,用户会希望推荐结果基于最近这几次评分进行一定的更
新,使得推荐结果匹配用户近期的偏好,满足用户近期的口味。
如果实时推荐继续采用离线推荐中的 ALS 算法,由于算法运行时间巨大,不
具有实时得到新的推荐结果的能力;并且由于算法本身的使用的是评分表,用户本
次评分后只更新了总评分表中的一项,使得算法运行后的推荐结果与用户本次评分
之前的推荐结果基本没有多少差别,从而给用户一种推荐结果一直没变化的感觉,
很影响用户体验。
另外,在实时推荐中由于时间性能上要满足实时或者准实时的要求,所以算法的计算量不能太大,避免复杂、过多的计算造成用户体验的下降。鉴于此,推荐精
度往往不会很高。实时推荐系统更关心推荐结果的动态变化能力,只要更新推荐结
果的理由合理即可,至于推荐的精度要求则可以适当放宽。
所以对于实时推荐算法,主要有两点需求:
(1)用户本次评分后、或最近几个评分后系统可以明显的更新推荐结果;
(2)计算量不大,满足响应时间上的实时或者准实时要求;
2、实时推荐算法设计
当用户 u 对电影 p 进行了评分,将触发一次对 u 的推荐结果的更新。由于用
户 u 对电影 p 评分,对于用户 u 来说,他与 p 最相似的电影们之间的推荐强度将
发生变化,所以选取与电影 p 最相似的 K 个电影作为候选电影。
每个候选电影按照“推荐优先级”这一权重作为衡量这个电影被推荐给用户 u
的优先级。
这些电影将根据用户 u 最近的若干评分计算出各自对用户 u 的推荐优先级,然
后与上次对用户 u 的实时推荐结果的进行基于推荐优先级的合并、替换得到更新后
的推荐结果。
具体来说:
首先,获取用户 u 按时间顺序最近的 K 个评分,记为 RK;获取电影 p 的最
相似的 K 个电影集合,记为 S;

sim(q,r)表示电影 q 与电影 r 的相似度,设定最小相似度为 0.6,当电影 q 和
电影 r 相似度低于 0.6 的阈值,则视为两者不相关并忽略;
sim_sum 表示 q 与 RK 中电影相似度大于最小阈值的个数;
incount 表示 RK 中与电影 q 相似的、且本身评分较高(>=3)的电影个数;
recount 表示 RK 中与电影 q 相似的、且本身评分较低(<3)的电影个数;
公式的意义如下:
首先对于每个候选电影 q,从 u 最近的 K 个评分中,找出与 q 相似度较高(>=0.6)的 u 已评分电影们,对于这些电影们中的每个电影 r,将 r 与 q 的相似 度乘以用户 u 对 r 的评分,将这些乘积计算平均数,作为用户 u 对电影 q 的评分 预测即

然后,将 u 最近的 K 个评分中与电影 q 相似的、且本身评分较高(>=3)的
电影个数记为 incount,计算 lgmax{incount,1}作为电影 q 的“增强因子”,意义
在于电影 q 与 u 的最近 K 个评分中的 n 个高评分(>=3)电影相似,则电影 q 的
优先级被增加 lgmax{incount,1}。如果电影 q 与 u 的最近 K 个评分中相似的高评
分电影越多,也就是说 n 越大,则电影 q 更应该被推荐,所以推荐优先级被增强
的幅度较大;如果电影 q 与 u 的最近 K 个评分中相似的高评分电影越少,也就是
n 越小,则推荐优先级被增强的幅度较小;
而后,将 u 最近的 K 个评分中与电影 q 相似的、且本身评分较低(<3)的电
影个数记为 recount,计算 lgmax{recount,1}作为电影 q 的“削弱因子”,意义在
于电影 q 与 u 的最近 K 个评分中的 n 个低评分(<3)电影相似,则电影 q 的优先
级被削减 lgmax{incount,1}。如果电影 q 与 u 的最近 K 个评分中相似的低评分电
影越多,也就是说 n 越大,则电影 q 更不应该被推荐,所以推荐优先级被减弱的
幅度较大;如果电影 q 与 u 的最近 K 个评分中相似的低评分电影越少,也就是 n 越
小,则推荐优先级被减弱的幅度较小;
最后,将增强因子增加到上述的预测评分中,并减去削弱因子,得到最终的 q 电
影对于 u 的推荐优先级。在计算完每个候选电影 q 的
后,将生成一组<电影 q 的 ID, q 的推荐优先级>的列表 updatedList:

而在本次为用户 u 实时推荐之前的上一次实时推荐结果 Rec 也是一组<电影
m,m 的推荐优先级>的列表,其大小也为 K:

接下来,将 updated_S 与本次为 u 实时推荐之前的上一次实时推荐结果 Rec进行基于合并、替换形成新的推荐结果 NewRec:

总之,实时推荐算法流程流程基本如下:
(1)用户 u 对电影 p 进行了评分,触发了实时推荐的一次计算;
(2)选出电影 p 最相似的 K 个电影作为集合 S;
(3)获取用户 u 最近时间内的 K 条评分,包含本次评分,作为集合 RK;
(4)计算电影的推荐优先级,产生<qID,>集合 updated_S;
将 updated_S 与上次对用户 u 的推荐结果 Rec 利用公式(4-4)进行合并,产生
新的推荐结果 NewRec;作为最终输出。
我们在 recommender 下新建子项目 StreamingRecommender,引入 spark、scala、
mongo、redis 和 kafka 的依赖:
MovieRecommendSystemrecommenderpom.xml
<dependencies>
<!-- Spark 的依赖引入 -->
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-core_2.11</artifactId>
</dependency>
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-sql_2.11</artifactId>
</dependency>
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-streaming_2.11</artifactId>
</dependency>
<!-- 引入 Scala -->
<dependency>
<groupId>org.scala-lang</groupId>
<artifactId>scala-library</artifactId>
</dependency>
<!-- 加入 MongoDB 的驱动 -->
<!-- 用于代码方式连接 MongoDB -->
<dependency>
<groupId>org.mongodb</groupId>
<artifactId>casbah-core_2.11</artifactId>
<version>2.8.0</version>
</dependency>
<!-- 用于 Spark 和 MongoDB 的对接 -->
<dependency>
<groupId>org.mongodb.spark</groupId>
<artifactId>mongo-spark-connector_2.11</artifactId>
<version>${mongodb-spark.version}</version>
</dependency>
<!-- redis -->
<dependency>
<groupId>redis.clients</groupId>
<artifactId>jedis</artifactId>
<version>2.9.0</version>
</dependency>
<!-- kafka -->
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>0.10.2.1</version>
</dependency>
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-streaming-kafka-0-10_2.11</artifactId>
<version>${spark.version}</version>
</dependency>
</dependencies>代码中首先定义样例类和一个连接助手对象(用于建立 redis 和 mongo 连接), 并在 StreamingRecommender 中定义一些常量:
在 resources 文件夹下引入 log4j.properties
MovieRecommendSystemrecommenderStreamingRecommendersrcmainresourceslog4j.properties
log4j.rootLogger=warn, stdout
log4j.appender.stdout=org.apache.log4j.ConsoleAppender
log4j.appender.stdout.layout=org.apache.log4j.PatternLayout
log4j.appender.stdout.layout.ConversionPattern=%d{yyyy-MM-dd HH:mm:ss,SSS} %5p --- [%50t] %-80c(line:%5L) : %m%n具体代码实现(实时推荐模块):
MovieRecommendSystemrecommenderStreamingRecommendersrcmainscalacomstudystatisticsStatisticsRecommender.scala
package com.study.statistics
import com.mongodb.casbah.commons.MongoDBObject
import com.mongodb.casbah.{MongoClient, MongoClientURI}
import org.apache.kafka.common.serialization.StringDeserializer
import org.apache.spark.SparkConf
import org.apache.spark.sql.SparkSession
import org.apache.spark.streaming.kafka010.{ConsumerStrategies, KafkaUtils, LocationStrategies}
import org.apache.spark.streaming.{Seconds, StreamingContext}
import redis.clients.jedis.Jedis
// 连接助手对象
object ConnHelper extends Serializable{
lazy val jedis = new Jedis("localhost")
lazy val mongoClient = MongoClient(MongoClientURI("mongodb://localhost:27017/recommender01"))
}
case class MongConfig(uri:String,db:String)
// 标准推荐
case class Recommendation(mid:Int, score:Double)
// 用户的推荐
case class UserRecs(uid:Int, recs:Seq[Recommendation])
//电影的相似度
case class MovieRecs(mid:Int, recs:Seq[Recommendation])
object StatisticsRecommender {
val MAX_USER_RATINGS_NUM = 20
val MAX_SIM_MOVIES_NUM = 20
val MONGODB_STREAM_RECS_COLLECTION = "StreamRecs"
val MONGODB_RATING_COLLECTION = "Rating"
val MONGODB_MOVIE_RECS_COLLECTION = "MovieRecs"
//入口方法
def main(args: Array[String]): Unit = {
val config = Map(
"spark.cores" -> "local[*]",
"mongo.uri" -> "mongodb://localhost:27017/recommender01",
"mongo.db" -> "recommender01",
"kafka.topic" -> "recommender01"
)
// 创建一个sparkConf
val sparkConf = new SparkConf().setMaster(config("spark.cores")).setAppName("SreamingRecommender")
// 创建一个SparkSession
val spark = SparkSession.builder().config(sparkConf).getOrCreate()
//拿到streaming context
val sc = spark.sparkContext
val ssc = new StreamingContext(sc, Seconds(2)) //batch duration
import spark.implicits._
implicit val mongConfig = MongConfig(config("mongo.uri"), config("mongo.db"))
//加载电影相似度矩阵数据,把它广播出去
val simMoviesMatrix = spark.read
.option("uri", mongConfig.uri)
.option("collection", MONGODB_MOVIE_RECS_COLLECTION)
.format("com.mongodb.spark.sql")
.load()
.as[MovieRecs]
.rdd
.map { moiveRecs => //为了查询相似度方便,转换成map
(moiveRecs.mid, moiveRecs.recs.map(x => (x.mid, x.score)).toMap)
}.collectAsMap()
//广播变量
val simMoviesMatrixBroadCast = sc.broadcast(simMoviesMatrix)
//实时流式处理
//创建到 Kafka 的连接
val kafkaParam = Map(
"bootstrap.servers" -> "spark2.x:9092",
"key.deserializer" -> classOf[StringDeserializer],
"value.deserializer" -> classOf[StringDeserializer],
"group.id" -> "recommender01",
"auto.offset.reset" -> "latest"
)
//通过kafka创建一个DStream
val kafkaStream = KafkaUtils.createDirectStream[String, String](
ssc, LocationStrategies.PreferConsistent,
ConsumerStrategies.Subscribe[String, String](Array(config("kafka.topic")), kafkaParam))
//把原始数据 UID|MID|SCORE|TIMESTAMP,(转化成)产生评分流
val ratingStream = kafkaStream.map {
msg =>
var attr = msg.value().split("|")
(attr(0).toInt, attr(1).toInt, attr(2).toDouble, attr(3).toInt)
}
//继续做流式处理,核心实时算法部分
ratingStream.foreachRDD {
rdds =>rdds.foreach {
case (uid, mid, score, timestamp) => {
println("rating data coming! >>>>>>>>>>>>>>>>")
//1、从redis里获取当前用户最近k次评分,保存成Array[(mid,score)]
val userRecentlyRatings = getUserRecentlyRating(MAX_USER_RATINGS_NUM, uid, ConnHelper.jedis)
//2、从相似度矩阵中取出当前电影最相似的N个电影,作为备选列表,Array[mid]
val candidateMovies = getTopSimMovies(MAX_SIM_MOVIES_NUM, mid, uid, simMoviesMatrixBroadCast.value)
//3、对每个备选电影,计算推荐优先级,得到当前用户的实时推荐列表,Array[(mid,score)]
val streamRecs = computeMovieScores(candidateMovies, userRecentlyRatings, simMoviesMatrixBroadCast.value)
//4、把推荐系统保存到mongodb
saveDataToMongoDB(uid, streamRecs)
}
}
}
//开始接收和处理数据
ssc.start()
println(">>>>>>>>>>>>> streaming started")
ssc.awaitTermination()
//redis操作返回的是java类,为了用map操作需要引入转化类
import scala.collection.JavaConversions._
def getUserRecentlyRating(num: Int, uid: Int, jedis: Jedis): Array[(Int, Double)] = {
//从用户的队列中取出 num 个评分数据保存在uid:UID为key的队列里,value是MID:SCORE
jedis.lrange("uid:" + uid.toString, 0, num)
.map {
item =>
val attr = item.split(":")
(attr(0).trim.toInt, attr(1).trim.toDouble)
}
.toArray
}
def getTopSimMovies(num: Int, mid: Int, uid: Int, simMovies: scala.collection.Map[Int, scala.collection.immutable.Map[Int, Double]])
(implicit mongConfig: MongConfig): Array[Int] = {
//1、从相似度矩阵中拿到所有相似的电影
val allSimMovies = simMovies(mid).toArray
//2、从mongodb中查询用户已看过的电影
val ratingExist = ConnHelper.mongoClient(mongConfig.db)(MONGODB_RATING_COLLECTION)
.find(MongoDBObject("uid" -> uid))
.toArray
.map {
item => item.get("mid").toString.toInt
}
//3、把看到的过滤,得到输出列表
allSimMovies.filter(x => ! ratingExist.contains(x._1))
.sortWith(_._2 > _._2)
.take(num)
.map(x => x._1)
}
def computeMovieScores(candidateMovies: Array[Int],
userRecentlyRatings: Array[(Int, Double)],
simMoives: scala.collection.Map[Int, scala.collection.immutable.Map[Int, Double]]): Array[(Int, Double)] = {
//定义一个ArrayBuffer,用于保存每个备选电影基础得分
val scores = scala.collection.mutable.ArrayBuffer[(Int, Double)]()
//定义一个HashMap,保存每一个备选电影增强减弱因子
val increMap = scala.collection.mutable.HashMap[Int, Int]()
val decreMap = scala.collection.mutable.HashMap[Int, Int]()
for (candiateMoive <- candidateMovies; userRecentlyRating <- userRecentlyRatings) {
//拿到备选电影和最近评分的相似度
val simScore = getMoviesSimScore(candiateMoive, userRecentlyRating._1, simMoives)
if (simScore > 0.7) {
//计算备选电影的基础推荐得分
scores += ((candiateMoive, simScore * userRecentlyRating._2))
if (userRecentlyRating._2 > 3) {
increMap(candiateMoive) = increMap.getOrDefault(candiateMoive, 0) + 1
} else {
decreMap(candiateMoive) = decreMap.getOrDefault(candiateMoive, 0) + 1
}
}
}
//根据备选电影的mid的groupby,根据公式去求最后的推荐评分
scores.groupBy(_._1).map {
//groupby 之后得到的数据 Map(mid ->)ArrayBuffer[(mid, score)])
case (mid, scoreList) =>
(mid, scoreList.map(_._2).sum / scoreList.length + log(increMap.getOrDefault(mid, 1)) - log(decreMap.getOrDefault(mid, 1)))
}.toArray
}
//获取两部电影的相似度
def getMoviesSimScore(mid1: Int, mid2: Int, simMoives: scala.collection.Map[Int,
scala.collection.immutable.Map[Int, Double]]) = {
simMoives.get(mid1) match {
case Some(sims) => sims.get(mid2) match {
case Some(score) => score
case None => 0.0
}
case None => 0.0
}
}
//求一个数的对数,底数默认10
def log(m: Int): Double = {
val N = 10
math.log(m) / math.log(N)
}
def saveDataToMongoDB(uid: Int, streamRecs: Array[(Int, Double)])(implicit mongConfig: MongConfig): Unit = {
//定义到StreamRescs表连接
val streamRecsCollection = ConnHelper.mongoClient(mongConfig.db)(MONGODB_STREAM_RECS_COLLECTION)
//如果表中已有uid对应的数据,则删除
streamRecsCollection.findAndRemove(MongoDBObject("uid" -> uid))
//将streamRecs数据存入表中
streamRecsCollection.insert(MongoDBObject("uid" -> uid,
"recs" -> streamRecs.map(x => MongoDBObject("mid" -> x._1, "score" -> x._2))))
}
}
}把zookeeper、kafka启动起来
[root@spark2 ~]# jps
10265 Kafka
10237 QuorumPeerMain
//查看当前kafka库的topic有哪些
[root@spark2 kafka]# ./bin/kafka-topics.sh --zookeeper spark2.x:2181 --list
__consumer_offsets
hotitems
recommender
sensor
sinkTest
//创建kafka生产者
[root@spark2 kafka]# ./bin/kafka-console-producer.sh --broker-list spark2.x:9092 --topic recommender01
>启动redis服务,进入redis客户端进行操作
//向redis里的uid库,插入电影的评分数据
127.0.0.1:6379> lpush uid:2 265,5.0 266,5.0 272,3.0 273,4.0 292,3.0 296,4.0 300,3.0
(integer) 7 //表示成功插入7条数据
//查看当前库,发现redis里多了一个uid库
127.0.0.1:6379> keys *
1) "1511679600000"
2) "uid:2"
//查看“uid”库下插入的详情信息(逆序存储)
127.0.0.1:6379> LRANGE uid:2 0 -1
1) "300,3.0"
2) "296,4.0"
3) "292,3.0"
4) "273,4.0"
5) "272,3.0"
6) "266,5.0"
7) "265,5.0"之后,就可以启动idea程序(StatisticsRecommender),启动之前,需要把mongodb服务启动,发现idea程序无报错:

我们可以在mongodb数据库查看信息:
> show dbs
admin 0.000GB
local 0.000GB
recommender 0.004GB
recommender01 0.081GB
> use recommender01
switched to db recommender01
> show tables
AverageMovies
GenresTopMovies
Movie
MovieRecs
RateMoreMovies
RateMoreRecentlyMovies
Rating
Tag
UserRecs
>二、基于内容的推荐服务
1、基于内容的推荐服务
原始数据中的 tag 文件,是用户给电影打上的标签,这部分内容想要直接转成
评分并不容易,不过我们可以将标签内容进行提取,得到电影的内容特征向量,进
而可以通过求取相似度矩阵。这部分可以与实时推荐系统直接对接,计算出与用户
当前评分电影的相似电影,实现基于内容的实时推荐。为了避免热门标签对特征提
取的影响,我们还可以通过 TF-IDF 算法对标签的权重进行调整,从而尽可能地接近
用户偏好, 然后通过电影特征向量进而求出相似度矩阵,就可以为实时推荐提供基础,得
到用户推荐列表了。可以看出,基于内容和基于隐语义模型,目的都是为了提取出
物品的特征向量,从而可以计算出相似度矩阵。而我们的实时推荐系统算法正是基
于相似度来定义的。
更新实时推荐结果
当计算出候选电影的推荐优先级的数组 updatedRecommends<movieId, E>后,这
个数组将被发送到 Web 后台服务器,与后台服务器上 userId 的上次实时推荐结果
recentRecommends<movieId, E>进行合并、替换并选出优先级 E 前 K 大的电影作为
本次新的实时推荐。具体而言:
a.合并:将 updatedRecommends 与 recentRecommends 并集合成为一个新的
<movieId, E>数组;
b.替换(去重):当 updatedRecommends 与 recentRecommends 有重复的电影
movieId 时,recentRecommends 中 movieId 的推荐优先级由于是上次实时推荐的结
果,于是将作废,被替换成代表了更新后的 updatedRecommends 的 movieId 的推荐
优先级;
c.选取 TopK:在合并、替换后的<movieId, E>数组上,根据每个 movie 的推荐
优先级,选择出前 K 大的电影,作为本次实时推荐的最终结果。
2、工程构建
在 recommender 下新建子项目 ContentRecommender,pom.xml 文件中只需引入 spark、scala 和 mongodb 的相关依赖:
MovieRecommendSystemrecommenderContentRecommenderpom.xml
<dependencies>
<dependency>
<groupId>org.scalanlp</groupId>
<artifactId>jblas</artifactId>
<version>${jblas.version}</version>
</dependency>
<!-- Spark的依赖引入 -->
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-core_2.11</artifactId>
</dependency>
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-sql_2.11</artifactId>
</dependency>
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-mllib_2.11</artifactId>
</dependency>
<!-- 引入Scala -->
<dependency>
<groupId>org.scala-lang</groupId>
<artifactId>scala-library</artifactId>
</dependency>
<!-- 加入MongoDB的驱动 -->
<dependency>
<groupId>org.mongodb</groupId>
<artifactId>casbah-core_2.11</artifactId>
<version>2.8.0</version>
</dependency>
<dependency>
<groupId>org.mongodb.spark</groupId>
<artifactId>mongo-spark-connector_2.11</artifactId>
<version>${mongodb-spark.version}</version>
</dependency>
</dependencies>在 resources 文件夹下引入 log4j.properties
MovieRecommendSystemrecommenderContentRecommendersrcmainresourceslog4j.properties
log4j.rootLogger=info, stdout
log4j.appender.stdout=org.apache.log4j.ConsoleAppender
log4j.appender.stdout.layout=org.apache.log4j.PatternLayout
log4j.appender.stdout.layout.ConversionPattern=%d{yyyy-MM-dd HH:mm:ss,SSS} %5p --- [%50t] %-80c(line:%5L) : %m%n具体代码实现(基于内容推荐模块):
MovieRecommendSystemrecommenderContentRecommendersrcmainscalacomstudycontentContentRecommender.scala
package com.study.content
import org.apache.spark.SparkConf
import org.apache.spark.ml.feature.{HashingTF, IDF, Tokenizer}
import org.apache.spark.ml.linalg.SparseVector
import org.apache.spark.sql.SparkSession
import org.jblas.DoubleMatrix
// 需要的数据源是电影内容信息
case class Movie(mid: Int, name: String, descri: String, timelong: String, issue: String,
shoot: String, language: String, genres: String, actors: String, directors: String)
case class MongoConfig(uri:String, db:String)
// 定义一个基准推荐对象
case class Recommendation( mid: Int, score: Double )
// 定义电影内容信息提取出的特征向量的电影相似度列表
case class MovieRecs( mid: Int, recs: Seq[Recommendation] )
object ContentRecommender {
// 定义表名和常量
val MONGODB_MOVIE_COLLECTION = "Movie"
val CONTENT_MOVIE_RECS = "ContentMovieRecs"
def main(args: Array[String]): Unit = {
val config = Map(
"spark.cores" -> "local[*]",
"mongo.uri" -> "mongodb://localhost:27017/recommender01",
"mongo.db" -> "recommender01"
)
val sparkConf = new SparkConf().setMaster(config("spark.cores")).setAppName("OfflineRecommender")
// 创建一个SparkSession
val spark = SparkSession.builder().config(sparkConf).getOrCreate()
import spark.implicits._
implicit val mongoConfig = MongoConfig(config("mongo.uri"), config("mongo.db"))
// 加载数据,并作预处理
val movieTagsDF = spark.read
.option("uri", mongoConfig.uri)
.option("collection", MONGODB_MOVIE_COLLECTION)
.format("com.mongodb.spark.sql")
.load()
.as[Movie]
.map(
// 提取mid,name,genres三项作为原始内容特征,分词器默认按照空格做分词
x => ( x.mid, x.name, x.genres.map(c=> if(c=='|') ' ' else c) )
)
.toDF("mid", "name", "genres")
.cache()
// 核心部分: 用TF-IDF从内容信息中提取电影特征向量
// 创建一个分词器,默认按空格分词
val tokenizer = new Tokenizer().setInputCol("genres").setOutputCol("words")
// 用分词器对原始数据做转换,生成新的一列words
val wordsData = tokenizer.transform(movieTagsDF)
// 引入HashingTF工具,可以把一个词语序列转化成对应的词频
val hashingTF = new HashingTF().setInputCol("words").setOutputCol("rawFeatures").setNumFeatures(50)
val featurizedData = hashingTF.transform(wordsData)
// 引入IDF工具,可以得到idf模型
val idf = new IDF().setInputCol("rawFeatures").setOutputCol("features")
// 训练idf模型,得到每个词的逆文档频率
val idfModel = idf.fit(featurizedData)
// 用模型对原数据进行处理,得到文档中每个词的tf-idf,作为新的特征向量
val rescaledData = idfModel.transform(featurizedData)
// rescaledData.show(truncate = false)
val movieFeatures = rescaledData.map(
row => ( row.getAs[Int]("mid"), row.getAs[SparseVector]("features").toArray )
)
.rdd
.map(
x => ( x._1, new DoubleMatrix(x._2) )
)
movieFeatures.collect().foreach(println)
// 对所有电影两两计算它们的相似度,先做笛卡尔积
val movieRecs = movieFeatures.cartesian(movieFeatures)
.filter{
// 把自己跟自己的配对过滤掉
case (a, b) => a._1 != b._1
}
.map{
case (a, b) => {
val simScore = this.consinSim(a._2, b._2)
( a._1, ( b._1, simScore ) )
}
}
.filter(_._2._2 > 0.6) // 过滤出相似度大于0.6的
.groupByKey()
.map{
case (mid, items) => MovieRecs( mid, items.toList.sortWith(_._2 > _._2).map(x => Recommendation(x._1, x._2)) )
}
.toDF()
movieRecs.write
.option("uri", mongoConfig.uri)
.option("collection", CONTENT_MOVIE_RECS)
.mode("overwrite")
.format("com.mongodb.spark.sql")
.save()
spark.stop()
}
// 求向量余弦相似度
def consinSim(movie1: DoubleMatrix, movie2: DoubleMatrix):Double ={
movie1.dot(movie2) / ( movie1.norm2() * movie2.norm2() )
}
}启动idea程序,代码正常启动运行,无报错:

我们可以再mongodb库查看数据,发现多了一张“ContentMovieRecs”,这张表作为“基于内容推荐模块的数据表”
> show tables
AverageMovies
ContentMovieRecs //内容推荐模块的数据表
GenresTopMovies
Movie
MovieRecs
RateMoreMovies
RateMoreRecentlyMovies
Rating
Tag
UserRecs
>查看这张ContentMovieRecs表数据,如图所示:

三、实时系统联调测试
我们的系统实时推荐的数据流向是:业务系统 -> 日志 -> flume 日志采集 ->
kafka streaming 数据清洗和预处理 -> spark streaming 流式计算。在我们完成实时推
荐服务的代码后,应该与其它工具进行联调测试,确保系统正常运行
1、工程构建
在 recommender 下新建子项目 KafkaStream,pom.xml 文件中只需引入 spark、scala 和 mongodb 的相关依赖:
MovieRecommendSystemrecommenderKafkaStreampom.xml
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<parent>
<artifactId>recommender</artifactId>
<groupId>org.study</groupId>
<version>1.0-SNAPSHOT</version>
</parent>
<modelVersion>4.0.0</modelVersion>
<artifactId>KafkaStream</artifactId>
<dependencies>
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-streams</artifactId>
<version>0.10.2.1</version>
</dependency>
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>0.10.2.1</version>
</dependency>
</dependencies> <build>
<finalName>kafkastream</finalName>
<plugins>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-assembly-plugin</artifactId>
<configuration>
<archive>
<manifest>
<mainClass>com.atguigu.kafkastream.Application</mainClass>
</manifest>
</archive>
<descriptorRefs>
<descriptorRef>jar-with-dependencies</descriptorRef>
</descriptorRefs>
</configuration>
<executions>
<execution>
<id>make-assembly</id>
<phase>package</phase>
<goals>
<goal>single</goal>
</goals>
</execution>
</executions>
</plugin>
</plugins>
</build>
</project>2、代码编写
MovieRecommendSystemrecommenderKafkaStreamsrcmainjavacomstudykafkastreamApplication.java
package com.study.kafkastream;
import org.apache.kafka.streams.KafkaStreams;
import org.apache.kafka.streams.StreamsConfig;
import org.apache.kafka.streams.processor.TopologyBuilder;
import java.util.Properties;
public class Application {
public static void main(String[] args) {
String brokers = "spark2.x:9092";
String zookeepers = "spark2.x:2181";
// 输入和输出的topic
String from = "log";
String to = "recommender01";
// 定义kafka streaming的配置
Properties settings = new Properties();
settings.put(StreamsConfig.APPLICATION_ID_CONFIG, "logFilter");
settings.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, brokers);
settings.put(StreamsConfig.ZOOKEEPER_CONNECT_CONFIG, zookeepers);
// 创建 kafka stream 配置对象
StreamsConfig config = new StreamsConfig(settings);
// 创建一个拓扑建构器
TopologyBuilder builder = new TopologyBuilder();
// 定义流处理的拓扑结构
builder.addSource("SOURCE", from)
.addProcessor("PROCESSOR", ()->new LogProcessor(), "SOURCE")
.addSink("SINK", to, "PROCESSOR");
KafkaStreams streams = new KafkaStreams( builder, config );
streams.start();
System.out.println("Kafka stream started!>>>>>>>>>>>");
}
}MovieRecommendSystemrecommenderKafkaStreamsrcmainjavacomstudykafkastreamLogProcessor.java
package com.study.kafkastream;
import org.apache.kafka.streams.processor.Processor;
import org.apache.kafka.streams.processor.ProcessorContext;
public class LogProcessor implements Processor<byte[], byte[]> {
//上下文
private ProcessorContext context;
@Override
public void init(ProcessorContext processorContext) {
this.context = processorContext;
}
@Override
public void process(byte[] dummy, byte[] line) {
// 把收集到的日志信息用string表示
String input = new String(line);
// 根据前缀MOVIE_RATING_PREFIX:从日志信息中提取评分数据
if( input.contains("MOVIE_RATING_PREFIX:") ){
System.out.println("movie rating data coming!>>>>>>>>>>>" + input);
input = input.split("MOVIE_RATING_PREFIX:")[1].trim();
context.forward( "logProcessor".getBytes(), input.getBytes() );
}
}
@Override
public void punctuate(long l) {
}
@Override
public void close() {
}
}
3、配置并启动 flume
在 flume 的 conf 目录下新建 log-kafka01.properties,对 flume 连接 kafka 做配置:
agent.sources = exectail
agent.channels = memoryChannel
agent.sinks = kafkasink
# For each one of the sources, the type is defined
agent.sources.exectail.type = exec
# 下面这个路径是需要收集日志的绝对路径,改为自己的日志目录
agent.sources.exectail.command = tail –f
/mnt/d/MovieRecommendSystem/businessServer/src/main/log/agent.log
agent.sources.exectail.interceptors=i1
agent.sources.exectail.interceptors.i1.type=regex_filter
# 定义日志过滤前缀的正则
agent.sources.exectail.interceptors.i1.regex=.+MOVIE_RATING_PREFIX.+
# The channel can be defined as follows.
agent.sources.exectail.channels = memoryChannel
# Each sink's type must be defined
agent.sinks.kafkasink.type = org.apache.flume.sink.kafka.KafkaSink
agent.sinks.kafkasink.kafka.topic = log
agent.sinks.kafkasink.kafka.bootstrap.servers = spark2.x:9092
agent.sinks.kafkasink.kafka.producer.acks = 1
agent.sinks.kafkasink.kafka.flumeBatchSize = 20
#Specify the channel the sink should use
agent.sinks.kafkasink.channel = memoryChannel
# Each channel's type is defined.
agent.channels.memoryChannel.type = memory
# Other config values specific to each type of channel(sink or source)
# can be defined as well
# In this case, it specifies the capacity of the memory channel
agent.channels.memoryChannel.capacity = 10000配置好后,启动 flume:
/bin/flume-ng agent -c ./conf/ -f ./conf/log-kafka01.properties -n agent-Dflume.root.logger=INFO,console4、前端对接
(1)把businessServer前端工程拷贝到MovieRecommendSystem的父工程下,如图:

接着,我们还需要把在子模块businessServer,引入到父项目MovieRecommendSystem进来

当成功引入进来后,businessServer模块变成正常项目,由上面的灰色变成深黑色

(2)启动Web服务程序


启动成功后:如图所示:

需要启动三个程序:StatisticsRecommender、Application(kafka过滤程序)、businessServer

访问url:http://localhost:8088/index.html


进入一个电影的主页面:


自己的github代码下载:https://github.com/Peng8868sky/MovieRecommendSyste
更多推荐
所有评论(0)