ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

基于Hadoop的电影推荐系统:从零搭建到离线召回实战

基于Hadoop的电影推荐系统:从零搭建到离线召回实战 简介这份资源是面向计算机相关专业在校学生、教师及企业员工的大数据课程设计参考包围绕Hadoop生态构建电影推荐系统适合作为小组大作业、毕业设计或项目初期立项演示的完整方案。压缩包共9个文件约100KB以xml配置文件、properties参数文件、classpath与project工程描述文件为主另含一个可运行的jar包覆盖Struts、Spring、c3p0连接池与log4j日志等典型Java Web与大数据整合配置便于快速导入IDE并理解工程结构。目前已有479人学习下载说明该方案在同类作业中具备一定参考价值。读者可从中获得一套经过测试运行成功的推荐系统实现思路包括Hadoop与Mahout协同调度的任务组织方式、依赖与数据源配置模板以及可在此基础上修改扩展的代码骨架用于课程设计、毕设答辩或功能二次开发帮助节省环境搭建与调试时间。1. 基于Hadoop的电影推荐系统从零搭建到离线召回的真实路径很多同学做大数据课程设计时第一反应是去网上找一份“电影推荐系统”的源码跑起来发现要么依赖缺失要么数据量小到根本用不上Hadoop最后只能截图交差。这个标题真正要解决的问题是如何用Hadoop生态完成一个能跑通、有数据流、能解释清楚每一步的离线推荐系统。它适合正在做大数据课程设计、毕业设计或者想理解“推荐系统在分布式环境下到底怎么落地”的开发者。核心不是算法多先进而是把数据采集、清洗、存储、计算、召回这条链路用Hadoop串起来。我见过太多小组把精力花在调参上结果连HDFS都没跑通这是典型的顺序错误。先让数据在集群里流动起来再谈推荐效果。2. 推荐系统的离线链路与Hadoop组件选型2.1 为什么离线推荐是课程设计的最优解推荐系统按实时性分为离线批处理、近线准实时和在线实时三类。课程设计通常只有两三周集群资源有限选离线批处理是最稳妥的。离线推荐的核心逻辑是用历史行为数据训练模型或计算相似度生成用户-物品推荐列表存入数据库供前端查询。它不要求毫秒级响应天然适合MapReduce和Spark批处理。从数据流看一个完整的离线推荐链路包括原始数据用户评分、电影元数据落到HDFS经过清洗和格式转换进入Hive数仓做特征工程再用MapReduce或Spark计算物品相似度或矩阵分解最后把Top-N结果导出到MySQL或HBase。这条链路里Hadoop承担的是存储和批计算底座Hive负责SQL化的数据加工Spark MLlib提供现成的协同过滤实现。选型时不要贪多把HDFSYARNHiveSpark跑通比堆一堆组件更有说服力。提示如果小组只有一台机器伪分布式足够完成全部流程不必强求多节点集群。2.2 组件版本与集群规划课程设计常见的坑是版本冲突。Hadoop 3.x和Spark 3.x搭配比较稳Hive用3.1.x可以兼容Hadoop 3.x。JDK必须用JDK 8JDK 11以上会遇到反射权限问题。下面是一个经过验证的版本组合组件版本作用JDK1.8.0_301基础运行环境Hadoop3.2.4HDFS YARNHive3.1.3数据仓库与SQL加工Spark3.2.1协同过滤计算MySQL5.7存储推荐结果Scala2.12.15Spark开发语言集群规划上伪分布式把NameNode和DataNode放同一台机器内存建议8GB以上。如果做三节点主节点跑NameNode和ResourceManager两个从节点跑DataNode和NodeManager。每台机器至少2核4GB否则YARN容器起不来。2.3 从HDFS到推荐结果的完整数据流数据流的设计决定了后续代码怎么写。我一般按这个顺序组织第一步原始数据上传到HDFS的/movie/data/raw/目录。MovieLens数据集是课程设计最常用的包含ratings.csv和movies.csvratings有userId、movieId、rating、timestamp四个字段。第二步用Hive建外部表指向HDFS原始目录做清洗过滤评分记录少于5条的用户、被评分少于10次的电影去掉时间戳异常的记录。清洗后的数据写入/movie/data/clean/。第三步用Spark读取清洗后的数据训练ALS矩阵分解模型或者用ItemCF计算物品相似度。ALS在Spark MLlib里是现成的几行代码就能跑适合课程设计。第四步把每个用户的Top-20推荐结果写入MySQL表user_recommend字段包括userId、movieId、score、rank。第五步前端用Flask或Spring Boot读MySQL展示。这一步不是Hadoop的重点但能让答辩时演示更完整。整个链路里HDFS负责原始数据和中间结果的存储YARN负责资源调度Hive做SQL化清洗Spark做模型计算。每个环节都有明确的输入输出目录方便排查问题。3. 伪分布式环境搭建与数据准备3.1 Hadoop伪分布式搭建的关键配置伪分布式搭建的教程网上很多但真正容易翻车的是配置文件里的几个参数。我按最小可用原则列出必须改的项。core-site.xml里指定HDFS的地址和临时目录configuration property namefs.defaultFS/name valuehdfs://localhost:9000/value /property property namehadoop.tmp.dir/name value/opt/module/hadoop-3.2.4/tmp/value /property /configurationhdfs-site.xml设置副本数为1因为伪分布式只有一个DataNodeconfiguration property namedfs.replication/name value1/value /property property namedfs.namenode.name.dir/name value/opt/module/hadoop-3.2.4/data/namenode/value /property property namedfs.datanode.data.dir/name value/opt/module/hadoop-3.2.4/data/datanode/value /property /configurationmapred-site.xml指定YARN作为资源调度器configuration property namemapreduce.framework.name/name valueyarn/value /property /configurationyarn-site.xml配置NodeManager的辅助服务configuration property nameyarn.nodemanager.aux-services/name valuemapreduce_shuffle/value /property property nameyarn.nodemanager.env-whitelist/name valueJAVA_HOME,HADOOP_COMMON_HOME,HADOOP_HDFS_HOME,HADOOP_CONF_DIR,CLASSPATH_PREPEND_DISTCACHE,HADOOP_YARN_HOME,HADOOP_MAPRED_HOME/value /property /configuration配置完成后格式化NameNodehdfs namenode -format。然后启动HDFS和YARNstart-dfs.sh和start-yarn.sh。用jps检查进程应该看到NameNode、DataNode、ResourceManager、NodeManager、SecondaryNameNode五个进程。注意格式化只能执行一次重复格式化会导致DataNode的clusterID和NameNode不一致DataNode起不来。如果误操作删掉data目录重新格式化。3.2 MovieLens数据上传与Hive外部表建立MovieLens数据集从GroupLens官网下载课程设计用ml-latest-small就够了约10万条评分。解压后得到ratings.csv和movies.csv。上传到HDFShdfs dfs -mkdir -p /movie/data/raw hdfs dfs -put ratings.csv /movie/data/raw/ hdfs dfs -put movies.csv /movie/data/raw/ hdfs dfs -ls /movie/data/raw/在Hive里建外部表指向这个目录CREATE EXTERNAL TABLE raw_ratings ( userId INT, movieId INT, rating FLOAT, ts BIGINT ) ROW FORMAT DELIMITED FIELDS TERMINATED BY , STORED AS TEXTFILE LOCATION /movie/data/raw/ratings.csv;这里有个细节ratings.csv第一行是表头Hive建表时用TBLPROPERTIES (skip.header.line.count1)跳过。如果数据里有引号包裹的字段需要加ESCAPED BY 。建好表后查一下数据量SELECT COUNT(*) FROM raw_ratings;。如果返回0检查文件路径和分隔符。3.3 用Hive SQL完成数据清洗与特征提取清洗的目标是去掉噪声数据生成用户-电影评分矩阵。我一般分三步第一步统计每个用户的评分数量和每个电影的被评分数量CREATE TABLE user_movie_count AS SELECT userId, COUNT(*) AS cnt FROM raw_ratings GROUP BY userId; CREATE TABLE movie_user_count AS SELECT movieId, COUNT(*) AS cnt FROM raw_ratings GROUP BY movieId;第二步过滤低频用户和低频电影保留评分记录CREATE TABLE clean_ratings AS SELECT r.userId, r.movieId, r.rating, r.ts FROM raw_ratings r JOIN user_movie_count u ON r.userId u.userId JOIN movie_user_count m ON r.movieId m.movieId WHERE u.cnt 5 AND m.cnt 10;第三步把清洗后的数据写入HDFS的clean目录INSERT OVERWRITE DIRECTORY /movie/data/clean ROW FORMAT DELIMITED FIELDS TERMINATED BY , SELECT userId, movieId, rating, ts FROM clean_ratings;这三步跑完数据量大概减少20%到30%但推荐质量会明显提升。低频用户和低频电影在协同过滤里贡献的噪声远大于信息量。参数说明u.cnt 5和m.cnt 10是经验阈值。用户评分少于5条行为模式不明显电影被评分少于10次相似度计算不可靠。如果数据集本身很小可以适当降低但不要低于3和5。4. 用Spark ALS跑通电影推荐召回4.1 ALS矩阵分解在推荐里的直观理解ALS叫交替最小二乘法是协同过滤里最常用的矩阵分解算法。它的核心思想是把用户-物品评分矩阵分解成两个低维矩阵一个表示用户隐向量一个表示物品隐向量。用户对物品的预测评分就是两个隐向量的点积。为什么叫“交替”因为同时优化两个矩阵很难ALS固定一个矩阵去优化另一个交替进行。Spark MLlib的ALS实现做了分布式优化能处理千万级评分数据。课程设计里用ml-latest-small单机跑几分钟就出结果。ALS有两个关键参数rank和regParam。rank是隐向量的维度越大模型容量越强但容易过拟合。regParam是正则化系数控制模型复杂度。课程设计里rank设10到50regParam设0.01到0.1效果比较稳。4.2 Spark读取HDFS数据并训练ALS模型用Scala写Spark程序或者用PySpark。课程设计里PySpark更友好代码量少。下面是一个完整的训练脚本from pyspark.sql import SparkSession from pyspark.ml.recommendation import ALS from pyspark.ml.evaluation import RegressionEvaluator # 初始化SparkSession指定master为yarn spark SparkSession.builder \ .appName(MovieRecommendALS) \ .master(yarn) \ .config(spark.executor.memory, 2g) \ .config(spark.executor.cores, 2) \ .getOrCreate() # 从HDFS读取清洗后的评分数据 ratings spark.read.csv(hdfs://localhost:9000/movie/data/clean, schemauserId INT, movieId INT, rating FLOAT, ts LONG) # 划分训练集和测试集 (training, test) ratings.randomSplit([0.8, 0.2], seed42) # 构建ALS模型 als ALS(maxIter10, rank20, regParam0.05, userColuserId, itemColmovieId, ratingColrating, coldStartStrategydrop, nonnegativeTrue) # 训练模型 model als.fit(training) # 在测试集上预测 predictions model.transform(test) # 计算RMSE evaluator RegressionEvaluator(metricNamermse, labelColrating, predictionColprediction) rmse evaluator.evaluate(predictions) print(RMSE str(rmse)) # 为每个用户生成Top-20推荐 userRecs model.recommendForAllUsers(20) # 结果写入HDFS userRecs.write.mode(overwrite).json(hdfs://localhost:9000/movie/data/recs) spark.stop()逻辑说明randomSplit按8:2划分训练和测试seed固定保证可复现。coldStartStrategydrop表示遇到训练集中没出现过的用户或物品时直接丢弃预测结果避免NaN。nonnegativeTrue强制隐向量非负对评分数据更合理。参数说明maxIter10是迭代次数一般5到20之间。rank20是隐向量维度数据量小可以降到10数据量大可以升到50。regParam0.05是正则化系数如果RMSE偏高可以调大到0.1如果欠拟合可以调小到0.01。4.3 推荐结果写入MySQL与前端查询Spark算出的推荐结果在HDFS上是JSON格式需要导出到MySQL供前端查询。用Spark的JDBC写入# 把userRecs展开成(userId, movieId, score)三列 from pyspark.sql.functions import explode flatRecs userRecs.select( userId, explode(recommendations).alias(rec) ).select( userId, rec.movieId, rec.rating ) # 写入MySQL flatRecs.write \ .format(jdbc) \ .option(url, jdbc:mysql://localhost:3306/movie_rec) \ .option(dbtable, user_recommend) \ .option(user, root) \ .option(password, your_password) \ .mode(overwrite) \ .save()MySQL表结构CREATE TABLE user_recommend ( userId INT, movieId INT, score FLOAT, PRIMARY KEY (userId, movieId) );前端用Flask写一个简单接口from flask import Flask, jsonify import pymysql app Flask(__name__) app.route(/recommend/int:user_id) def recommend(user_id): conn pymysql.connect(hostlocalhost, userroot, passwordyour_password, dbmovie_rec) cursor conn.cursor() cursor.execute( SELECT movieId, score FROM user_recommend WHERE userId%s ORDER BY score DESC LIMIT 20, (user_id,)) results cursor.fetchall() conn.close() return jsonify([{movieId: r[0], score: r[1]} for r in results]) if __name__ __main__: app.run(host0.0.0.0, port5000)这样从HDFS到MySQL再到前端展示的链路就通了。答辩时打开浏览器输入http://localhost:5000/recommend/1能看到推荐列表比只贴代码截图有说服力。5. 避坑与排查课程设计里最容易翻车的五个点5.1 DataNode起不来jps看不到进程现象执行start-dfs.sh后jps只显示NameNode和SecondaryNameNode没有DataNode。原因最常见的是重复格式化NameNode导致DataNode的clusterID和NameNode不一致。其次是data目录权限不对或者磁盘空间不足。解决先看DataNode日志logs/hadoop-*-datanode-*.log如果报Incompatible clusterIDs删掉data目录下所有内容重新格式化。如果是权限问题用chown -R hadoop:hadoop /opt/module/hadoop-3.2.4/data。格式化前确认没有残留进程。5.2 Hive查询报错“FAILED: SemanticException Unable to instantiate”现象Hive建表或查询时报语义异常提示无法实例化某个类。原因Hive的元数据库没初始化或者MySQL驱动没放到Hive的lib目录。解决先检查hive-site.xml里JDBC连接串、用户名、密码是否正确。然后执行schematool -dbType mysql -initSchema初始化元数据库。最后确认mysql-connector-java-5.1.xx.jar在$HIVE_HOME/lib下。如果还报错看/tmp/hive目录权限确保当前用户可写。5.3 Spark提交到YARN后卡在ACCEPTED状态现象spark-submit提交后YARN的Web UI显示应用一直处于ACCEPTED不进入RUNNING。原因YARN资源不足或者队列配置限制了资源。伪分布式下如果yarn.nodemanager.resource.memory-mb设得太小容器申请不到内存。解决在yarn-site.xml里把yarn.nodemanager.resource.memory-mb调到4096以上yarn.scheduler.maximum-allocation-mb也调到4096。然后重启YARN。提交Spark时把spark.executor.memory降到1g或2g避免单个容器申请过大。5.4 ALS训练报NaN或RMSE异常高现象模型训练完预测结果里出现NaN或者RMSE超过2.0。原因训练集里有用户或电影在测试集里没出现过冷启动导致预测失败。或者评分数据没有做归一化ALS对量纲敏感。解决设置coldStartStrategydrop丢弃冷启动预测。检查评分范围MovieLens是0.5到5.0如果数据里有0分或负分先过滤掉。另外regParam太小会导致过拟合适当调大到0.1试试。5.5 推荐结果全部一样没有个性化现象给不同用户生成的Top-20推荐列表几乎相同。原因数据量太小或者ALS的rank设得太低模型学不到用户差异。也可能是热门电影占据了所有推荐位。解决先检查清洗后的数据量如果评分记录少于1万条ALS很难学到个性化。可以换用ItemCF基于物品相似度做推荐对小数据集更友好。另外把rank从10提高到30或50增加模型容量。如果还是不行在推荐结果里加入多样性策略比如每个用户强制推荐至少3个不同类别的电影。6. 从离线召回走向可解释推荐一个课程设计的加分技巧课程设计答辩时老师最常问的问题是“为什么给这个用户推荐这部电影”如果只能回答“模型算出来的”分数不会高。我一般会加一个可解释性模块用物品相似度解释推荐理由。具体做法是在Spark里计算物品-物品相似度矩阵用余弦相似度或Jaccard相似度。然后对每个推荐结果找到用户历史评分最高的电影输出“因为你喜欢《肖申克的救赎》所以推荐《教父》”。这个逻辑用Spark SQL就能实现# 计算物品相似度 from pyspark.ml.feature import VectorAssembler from pyspark.ml.linalg import Vectors # 构建用户-物品评分矩阵 user_item ratings.groupBy(userId).pivot(movieId).avg(rating).fillna(0) # 转成向量 assembler VectorAssembler(inputColsuser_item.columns[1:], outputColfeatures) item_vectors assembler.transform(user_item) # 用余弦相似度计算物品相似度 from pyspark.mllib.linalg import Vectors as MLLibVectors from pyspark.mllib.linalg.distributed import RowMatrix rows item_vectors.select(features).rdd.map(lambda x: MLLibVectors.dense(x[0])) mat RowMatrix(rows) sim mat.columnSimilarities()这段代码算出物品之间的相似度然后对每个推荐结果找到用户评分最高的物品输出相似度最高的推荐物品作为解释。答辩时展示“推荐理由”列比单纯列推荐列表更有说服力。另一个加分点是加一个简单的评估指标。除了RMSE还可以算PrecisionK和RecallK。PrecisionK表示推荐列表前K个里有多少是用户实际喜欢的RecallK表示用户实际喜欢的物品有多少被推荐出来了。课程设计里K取10或20用测试集算一下写在报告里。我自己的习惯是先把离线链路跑通再加可解释性和评估指标。不要一开始就追求算法复杂度把HDFS、Hive、Spark、MySQL这条链路走通比换更高级的模型更有价值。课程设计考察的是工程能力不是算法创新。希望帮到你。本文还有配套的精品资源点击获取
返回列表