文章目录
- 每日一句正能量
- 第八章 Spark MLlib 机器学习算法库
- 章节概要
- 8.5 分类
- 8.5.1 线性支持向量机
- 8.5.2 逻辑回归
- 8.6 案例——构建推荐系统
- 8.6.1 推荐模型分类
- 8.6.2利用MLlib实现电影推荐
每日一句正能量
“用平和之心过朝朝暮暮,用知足之心览人间烟火。”
平和与知足是生活的最佳状态。朝朝暮暮的平凡,因心境而变得珍贵。
慢下来,在平凡中体味生活的丰盛,用知足平和的心态过好每一天。用这颗修炼好的心,去点染和照亮每一个平凡的日子,最终将整个生活过成一件艺术品。
第八章 Spark MLlib 机器学习算法库
章节概要
MLlib是Spark提供的处理机器学习方面的功能库,该库包含了许多机器学习算法,开发者可以不需要深入了解机器学习算法就能开发出相关程序。本章将介绍Spark MLlib基本知识以及使用方法,最后通过构建推荐引擎了解机器学习系统的构建思路及流程。
8.5 分类
MLlib支持多种分类分析方法,例如二元分类、多元分类,表列出了不同种类的问题可以采用不同的分类算法。
| 分析方法 | 相关算法 |
|---|---|
| 二元分类 | 线性支持向量机、逻辑回归、决策树、随机森林、梯度提升树、朴素贝叶斯 |
| 多元分类 | 逻辑回归、决策树、随机森林、朴素贝叶斯 |
分类是指将事物分成不同类别,在分类模型中,可根据一组特征来判断类别,这些特征代表了物体、事物或上下文的相关属性。分类算法又被称为分类器,它是数据挖掘和机器学习领域中的一个重要分支,它属于有监督学习的一种形式,我们用带有类标记或者类输出的训练样本来训练模型,要想评价一个分类器的好坏,我们就要有评价指标,最常见的就是准确率。
8.5.1 线性支持向量机
线性支持向量机是一种常见判别方法,在机器学习领域中是一个有监督学习模型,用来进行模式识别、分类以及回归分析。使用MLlib提供的线性支持向量机算法训练模型,需要导入线性支持向量机所需包。
#导入线性支持向量机所需包 scala>importorg.apache.spark.mllib.classification.{SVMModel,SVMWithSGD}#导入二元分类评估类 scala>importorg.apache.spark.mllib.evaluation.BinaryClassificationMetrics #MLUtils提供了一些辅助方法,用于加载,保存和预处理MLlib中使用的数据 scala>importorg.apache.spark.mllib.util.MLUtils #加载Spark官方提供数据集 scala>valdata=MLUtils.loadLibSVMFile(sc,"file:///export/servers/spark/data/mllib/sample_libsvm_data.txt")#将数据的60%分为训练数据,40%分为测试数据 scala>valsplits=data.randomSplit(Array(0.6,0.4),seed=11L)scala>valtraining=splits(0).cache()scala>valtest=splits(1)#设置迭代次数 scala>valnumIterations=100#执行算法来构建模型 scala>valmodel=SVMWithSGD.train(training,numIterations)#用测试数据评估模型 scala>valscoreAndLabels=test.map{point=>valscore=model.predice(point.features)(score,point.label)}#获取评估指标 scala>valmetrics=newBinaryClassificationMetrics(scoreAndLabels)#计算二元分类的PR和ROC曲线下的面积 scala>valauROC=metrics.areaUnderROC()auROC:Double=1.0#保存并加载模型 scala>model.save(sc,"target/tmp/scalaSVMWithSGDModel")scala>valsameModel=SVMModel.load(sc,"target/tmp/scalaSVMWithSGDModel")上述代码中,我们将数据文件分为两份,其中60%的数据为训练模型数据,40%的数据为测试数据,用来评估我们创建的模型。第19行代码调用SVMWithSGD.train()方法构建训练模型。为了检验分类器的好坏程度,可以利用MLlib提供的二元分类评估类计算ROC面积,ROC曲线是对分类器的真假阳性率图形化的解释,ROC下的面积(通常称为AUC)表示平均值,当AUC为1.0时,表示是一个完美的分类器,当AUC为0.5时,表示该模型和随机预测效果一样,没有必要使用。
评估模型完成后,还可以使用save()方法保存至HDFS目录中,下次只需调用load()方法即可加载该模型。结果如下图所示:
8.5.2 逻辑回归
逻辑回归又称为逻辑回归分析,是一个概率模型的分类算法,用于数据挖掘、疾病自动诊断及经济预测等领域。例如在流行病学研究中,探索引发某一疾病的危险因素,根据模型预测在不同自变量情况下,推测发生某一疾病。
Spark MLlib提供了逻辑回归算法,下面具体演示加载数据并执行训练模型方法,具体代码如下。
#导入逻辑回归所需包 scala>importorg.apache.spark.mllib.classification.{LogisticRegressionModel,LogisticRegressionWithLBFGS}#导入分类评估器 scala>importorg.apache.spark.mllib.evaluation.MulticlassMetrics scala>importorg.apache.spark.mllib.regressin.LabeledPoint scala>importorg.apache.spark.mllib.util.MLUtils #加载Spark官方提供数据集 scala>valdata=MLUtils.loadLibSVMFile(sc,"file:///export/servers/spark/data/mllib/sample_libsvm_data.txt")#将数据的60%分为训练数据,40%分为测试数据 scala>valsplits=data.randomSplit(Array(0.6,0.4),seed=11L)scala>valtraining=splits(0).cache()scala>valtest=splits(1)#运行训练算法来构建模型 scala>valmodel=newLogisticRegressionWithLBFGS().setNumClasses(10).run(training)#用测试数据评估模型 scala>valpredictionAndLabels=test.map{caseLabeledPoint(label,features)=>valprediction=model.predict(features)(prediction,label)}#获取评估指标 scala>valmetrics=newMulticlassMetrics(predictionAndLabels)scala>valaccuracy=metrics.accuracy accuracy:Double=1.0#保存并加载模型 scala>model.save(sc,"target/tmp/scalaLogisticRegressionWithLBFGSModel")scala>valsameModel=LogisticRegressionModel.load(sc,"target/tmp/scalaLogisticRegressionWithLBFGSModel")评估模型的性能不仅仅只有通过ROC曲线,通常在二元分类中使用的评估方法有:预测正确率和错误率、准确率和召回率等。准确率通常用于评估结果的质量,召回率用来评估结果的完整性,在二元分类问题中,准确率定义为真阳性数据个数除以真阳性和假阳性的数据总数,其中真阳性是指被正确预测的类别为1的样本,假阳性是错误预测为类别1的样本。如果每个数据被分类器预测为1的样本,那么准确率即为1.0。结果如下图所示:
8.6 案例——构建推荐系统
8.6.1 推荐模型分类
随着电子商务规模的不断扩大,商品个数和种类快速增长,顾客就需要花费大量的时间才能找到自己
想买的商品,这样就会造成消费者花费很长时间搜索商品,从而造成用户体验下降。为了解决这些问
题,个性化推荐系统应运而生。个性化推荐系统是建立在海量数据挖掘基础上的一种高级商务智能平
台,从而为其顾客购物提供完全个性化的决策支持和信息服务。
推荐系统的研究已经相当广泛,也是最为大众所知的一种机器学习模型,目前最为流行的推荐系统所
应用的算法是协同过滤,协同过滤通常用于推荐系统,这项技术是为了填补关联矩阵的缺失项,从而
实现推荐效果。简单地说,协同过滤是利用大量已有的用户偏好,来估计用户对其未接触的物品的喜好程度。
在协同过滤算法中有着两个分支:基于群体用户的协同过滤(UserCF)和基于物品的协同过滤(ltemCF).
1.基于物品的推荐(ltemCF)
基于物品的推荐是利用现有用户对物品的偏好或是评级情况,计算物品之间的某种相似度,以用户接触过的物品来
表示这个用户,然后寻找出和这些物品相似的物品,并将这些物品推荐给用户。
2.基于物品的推荐(UserCF)
基于用户的推荐,可以用“志趣相投”一词所表示,通常是对用户的历史行为数据分析,例如购买、收藏的商品,评论内容或搜索内容,
通过某种算法将用户喜好的物品进行打分。根据不同用户对相同物品或内容数据的态度和偏好程度
来计算用户之间的关系程度,在有相同喜好的用户之间进行商品推荐。
8.6.2利用MLlib实现电影推荐
在电影推荐系统中,通常分类针对用户推荐电影和针对电影推荐用户两种方式。具体实现方式取决于采用的推荐模型,若采用基于用户的推荐模型,则会利用相似用户的评级来计算对某个用户的推荐。若采用基于物品的推荐模型,则会依靠用户接触过的物品与候选物品之间的相似度来获得推荐。
在Spark MLlib实现了交替最小二乘(ALS)算法,它是机器学习的协同过滤式推荐算法,机器学习的协同过滤式推荐算法是通过观察所有用户给产品的评分来推断每个用户的喜好,并向用户推荐合适的产品。
接下来我们将分步骤讲解,利用Spark MLlib实现电影推荐案例的核心过程。
- 准备训练模型数据
MovieLens是历史最悠久的推荐系统,它是由美国Minnesota大学计算机科学与工程学院的GroupLens项目组创办,是一个以研究为目的的、非商业性质的实验性站点,读者可以从该网站中下载实验数据进行学习,网站地址为:https://grouplens.org/datasets/movielens/,下载ml-100k.zip解压包。具体如图8-6所示。
图8-6
也可以直接在Linux系统上输入以下命令下载文件,具体命令如下。
wgethttp://files.grouplens.org/datasets/movielens/m1-100k.zip实验数据文件下载完成后,将其进行解压,命令如下。
yuminstallunzipunzip-jm1-100k结果如下图所示:
最终将解压文件上传到HDFS中的/spark/mldata路径下,效果如图8-7所示。命令如下:
>hadoop fs-mkdir/spark/mldata>cd..>hadoop fs-putml-100k /spark/mldata结果如下图所示:
在本案例中,主要用到u.data文件(用户评分数据)以及u.item文件(电影数据),数据片段分别为如图所示。
文件u.item中,具有多个字段,本案例主要使用第一列电影id、第二列电影名称,后续将针对该文本进行字符串处理。
- 编写程序,训练模型
我们采用Spark-Shell读取u.data数据文件,将其转换为RDD,执行命令如下。
$ spark-shell--master local[2]#读取文件转换RDD scala>valdataRdd=sc.textFile("/spark/mldata/ml-100k/u.data")#输出RDD第一行数据 scala>dataRdd.first()res0:String=1922423881250949结果如下图所示:
从上一章节已经得知,该数据是由用户id、电影id、等级评价和时间戳依次组成,在训练模型时,可以去除时间戳字段,使用take()方法提取前三个字段即可,具体代码如下。
scala>valdataRdds=dataRdd.map(_.split("\t").take(3))scala>dataRdds.first()res1:Array[String]=Array(196,242,3)根据图的u.data文件内容可知,使用“\t”进行分割,返回一个Array[String]类型的RDD,分别对应用户id、影片id以及等级。至此就有了dataRdds数据集,可以使用first()函数查看第一行数据。
下面就可以使用Spark MLlib训练模型了,首先导入MLlib实现的ALS算法模型库。
scala>importorg.apache.spark.mllib.recommendation.ALS在ALS库中,可以通过调用train()函数来训练模型,具体代码如下。
def train( ratings: RDD[Rating], rank: Int, iterations: Int, lambda: Double ): MatrixFactorizationModel上述代码中,train()函数需要提供四个参数,如表所示。
| 参数名称 | 相关说明 |
|---|---|
| ratings | 训练的数据格式是Rating(UserID, productID, rating)的RDD |
| rank | 对应ALS模型中的因子个数,也就是在低阶近似矩阵中的隐含特征个数,因子个数一般越多越好,但是也会加大内存开销开销,通常值为10-200 |
| iterations | 对应运算时的迭代次数,减少评级矩阵的重建误差,默认值为5,大部分情况下设置10次左右 |
| lambda | 该参数控制模型的正则化过程,从而控制模型的拟合程度。值越高,正则化越严厉,该参数的值与实际数据的大小、特征和稀疏程度有关,默认值0.01 |
训练模型需要Rating格式的数据,可以将dataRdds使用map()方法进行转换,得到Rating格式数据,传入到train()函数,具体代码如下。
#导入Rating包 scala>importorg.apache.spark.mllib.recommendation.Rating scala>valratings=dataRdds.map{caseArray(user,movie,rating)=>Rating(user.toInt,movie.toInt,rating.toDouble)}scala>ratings.first()res6:org.apache.spark.mllib.recommendation.Rating=Rating(196,242,3.0)结果如下图所示:
需要注意的是,使用case语句来提供各属性对应的变量名,dataRdds是从u.data文本文件中转换的数据,因此需要把String类型转换成对应的数据类型,提取简单特征后,就可以调用train()函数训练模型,代码如下。
scala>valmodel=ALS.train(ratings,50,10,0.01)model:org.apache.spark.mllib.recommendation.MatrixFactorizationModel=org.apache.spark.mllib.recommendation.MatrixFactorizationModel@6580f76c结果如下图所示:
调用ALS.train训练数据集后,就会创建推荐引擎模型MatrixFactorizationModel(矩阵分解)对象,该对象成员如表所示。
| 对象成员 | 相关说明 |
|---|---|
| predict(user: Int, product: Int): Double | 计算给定用户和物品的预期得分 |
| productFeatures:RDD[(Int, Array[Double])] | 分解后的物品矩阵 |
| rank: Int | 分解后的参数 |
| userFeatures: Rdd[(Int, Array[Double])] | 分解后产品矩阵 |
表中,predict函数以(user, product)作为输入参数,该函数将为每一对生成相应的预测得分,具体代码如下。
scala>valpredictedRating=model.predict(100,200)predictedRating:Double=1.1136730131397399结果如下图所示:
从上述执行结果可以看出,改模型预测用户id=100对电影id=200的评级约为1.11。需要注意的是,ALS模型的初始化过程根据硬件环境以及参数等因素会造成不同的结果。
- 为用户推荐多个电影
如果要为某个用户推荐多个物品,可以调用MatrixFactorizationModel对象所提供的recommendProducts(user: Int, num: Int)函数来实现,返回值即为预测得分最高的前num个物品,,具体代码如下。
#定义用户id scala>valuserid=100#定义推荐数量 scala>valnum=10scala>valtopRecoPro=model.recommendProducts(userid,num)结果如下图所示:
从上述代码可以看出,使用训练完成的模型进行推荐,传入参数(user=100, num=10),返回结果是一个Rating数据类型的数组,其中参数分别表示用户id(user)、推荐电影id(product)、算法得出的评分(rating),其中评分越高,代表推荐引擎优先推荐这件物品。Rating(100, 207, 5.704436943409341)数据表示针对id=100的用户,预测对id=207的电影,评级为5.70分。
为了更直观的检测推荐效果,可以将u.item文件中的电影id与电影名称进行映射,因此首先读取u.item文件并转换为RDD。具体代码如下。
scala>valmoviesRdd=sc.textFile("/spark/mldata/ml-100k/u.item")根据图中u.item文件的数据格式进行分析,可以通过"|"字符分割,使用map()函数针对每一项数据进行转换,提取前2个数据,并将电影id、电影名称产生映射关系。具体代码如下。
valtitles=moviesRdd.map(line=>line.split("\\|").take(2)).map(array=>(array(0).toInt,array(1))).collectAsMap()结果如下图所示:
对于100个用户,可以通过Rating对象的rating属性来对推荐的电影名称进行匹配,具体代码如下。
scala>topRecoPro.map(rating=>(titles(rating.product),rating.rating)).foreach(println)结果如下图所示:
至此,根据用户推荐电影实现完成。
- 将物品推荐给用户
如果要为某个物品推荐多个用户时,可以调用MatrixFactorizationModel对象所提供的recommendUsers(product: Int, num: Int)函数来实现,其中product参数代表被推荐的物品id,num为推荐物品的数量,最终返回值为针对这件物品可能感兴趣的num名用户,具体代码如下。
scala>model.recommendUsers(100,5)通过上述返回结果看出,电影编号为100的推荐给用户编号为495、30、272、8、68这五位用户,至此,基于物品推荐电影实现完成。
转载自:https://blog.csdn.net/u014727709/article/details/133902769
欢迎 👍点赞✍评论⭐收藏,欢迎指正