ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

基于Spark的用户画像电影推荐系统设计实现

基于Spark的用户画像电影推荐系统设计实现 简介一份基于 Spark 的用户画像电影推荐系统设计项目资料包面向 Python、大数据方向的毕业设计或课程设计使用者重点解决用户行为数据清洗、用户画像标签构建、协同过滤推荐算法实现及结果展示的全流程问题。压缩包共 798 个文件大小约 15.54MB文件类型以 Python 源码py/pyc和前端页面资源html/css/js为主辅以 SQL 数据库脚本、PDF/文本说明文档及常见静态资源便于按模块对照查看后端算法、前端交互与数据表结构。目前已有 33 人浏览学习。包内 README 详细说明了系统架构、安装配置步骤与使用方法整体采用模块化设计涵盖 Spark MLlib 机器学习调用、协同过滤推荐、MySQL 存储、服务端请求响应等关键环节同时体现系统的可扩展性与健壮性设计能够帮助读者快速理解从数据到推荐的完整链路。对于完成毕业设计或课程设计的同学这份资料既可作为直接参考框架也可基于源码二次开发并支持撰写论文或设计文档时引用相关流程图与代码结构。1. 为什么说“画像Spark”是电影推荐系统里最容易做崩的一环拿到“基于Spark的用户画像电影推荐系统设计”这个标题时我先看到的不是 Spark 三个字而是“用户画像”这半截。绝大多数课程设计、毕设甚至小公司线上的推荐工程翻车都不是翻在算法上而是翻在画像标签没有形成能用的资产——所谓“画像”只是一个写死的用户表推荐却用的是另一套完全不相干的规则两头各跑各的模型上线之后就是黑匣子。这套设计真正该解决的是一件事把用户行为聚合成标签资产再用同一套标签资产去驱动召回、过滤和打分。所以它不是“装个 Spark 跑个 ALS”这么简单而是要把数据清洗、标签加工、协同过滤召回、候选过滤和业务打分串成一条完整链路。适合谁准备做课程设计、需要把 Scala 工程能力落地的数据开发以及正在把“用户画像”从 PPT 落到代码清单的技术负责人。下面我用一个可复现的方案把整条链路拆开讲。2. 把“用户画像”拆成标签生产链路四类标签与两个时空口径2.1 四类标签怎么分工事实、统计、规则、挖掘各管一段电影场景下的用户标签我一般会拆成四个层级来建而不是混在一张宽表里一把梭。第一层是事实标签也叫原始标签。它几乎不做加工直接记录“用户 U 在某时刻对电影 M 产生了行为”比如{uid, mid, action: play, ts, duration, device}。事实标签的价值在于可回溯、可核对也是后三层计算的输入。设计里我通常用 Parquet 按日期分区落一张事件明细表不更新、只追加。第二层是统计标签用聚合算出短周期内的行为画像例如“近 7 天观影 12 部”“平均完播率 0.6”“偏好电影类型 TOP3”。这层标签是推荐召回阶段性价比最高的信号因为它稳定、不容易被误解而且只要有行为就能算。统计标签的时效口径要特别定义下面小节单独讲。第三层是规则标签把业务经验编码成 if-else。举几个例子连续 7 天没有打开的用户打上“潜在流失”观影时长 top 20% 的用户打上“深度观影”只标记想看的用户打上“种草型用户”。规则标签写起来很快但维护最麻烦因为业务一变规则就要跟着改。我在项目里会把规则配置抽成 Scala 的 Map 结构而不是散落在 if 语句里。第四层是挖掘标签用模型产出的结果回填到画像中。典型代表是内容偏好向量把 ALS 隐因子或内容 Embedding 输出成{uid, feature_vector}表这个表既能为“相似用户”类召回提供信号也能直接写到用户画像里当作高维偏好标签。2.2 标签时效口径离线 T1 和实时近场画像的边界画像时效搞不清楚后面推荐评估会收到一堆“玄学波动”。我用的设计是双口径离线画像走 T1 全量重算近场画像走窗口实时聚合。离线画像跑在 Spark 批任务上凌晨对前一天的用户行为做全量聚合产出user_profile_offline表服务端把标签缓存下来推荐服务直接读缓存。近场画像则只保留用户最近 2 小时的行为摘要存放最近看过的电影、最近一次完播率、最近一次搜索词由 Flink 或 Spark Streaming 的小时级微批写入 Redis/HBase。上线时容易搞混的是实时画像和离线画像的数据口径不一致导致推荐服务里同一个用户两个画像来回覆盖。我的防线是固定两个口径——离线画像服务“中长期偏好”实时画像只服务“本次会话的近期兴趣”两者不互相覆盖离线画像尽力而为实时画像有窗口期。2.3 画像标签落到 Spark 上的基本模型事件→清洗→聚合→输出最稳的写法是先落地一份“画像聚合逻辑”的伪代码设计再写 Spark 代码。因为整个项目的接收者往往既要看文档又要能跑伪代码能先把聚合口径固定下来。// 画像聚合核心逻辑Scala 伪码 import org.apache.spark.sql.functions._ case class ActionEvent(uid: String, mid: String, action: String, ts: Long, duration: Long, score: Double) val eventDF spark.read.parquet(/data/events/dt2025-01-01) .as[ActionEvent] .filter($action.isin(play, collect, search_click, review)) .filter($duration 10) // 过滤掉误触 val profileDF eventDF .groupBy(uid) .agg( count(*).as(total_actions), sum(when($action play, $duration)).as(total_play_duration), avg(when($action play, $duration)).as(avg_play_duration), countDistinct(mid).as(distinct_movies) ) .withColumn(preference_vector, collect_list($mid).over(Window.partitionBy(uid).orderBy($ts)))逻辑说明先按dt读到对应日期分区的事件表过滤掉无效动作再按uid聚合出用户的统计画像Window部分用来输出最近观影序列这是之后 ALS 训练矩阵的原料。注意这里我临时用了preference_vector的窗口写法实际上全量跑的时候窗口开销很大生产环境我会改成把“最近 30 个 mid 序列”落成一个数组字段而不是在聚合里实时计算。参数说明duration 10是过滤误触的最小观影时长阈值课设里可以用 5 秒线上我见过 30 秒的保守版本total_actions和distinct_movies是最关键的两个统计字段后面做“活跃度分层”要靠它们。隐藏坑是avg(when(...))里如果不写otherwise(null)when为 false 时返回 nullavg 会自动忽略 null但如果你提前把 duration 转成了 0avg 就被拉低了。我见过一次因为 0 值导致“用户平均观影时长全低于 1 分钟”的翻车案例原因就是没处理缺失值。3. 从画像到推荐ALS 召回、候选过滤与业务打分三层架构3.1 为什么用 ALS对隐式反馈友好且 Spark MLlib 原生实现基于 Spark 的电影推荐系统召回算法我首选 ALS交替最小二乘理由有三条。第一它是 Spark MLlib 自带的实现不需要额外对接外部算法库从工程落地角度风险最小。我的实测经验是只要你数据量在千万级以下spark.ml.recommendation.ALS能稳定跑完不会因为矩阵太大而 OOM。第二ALS 对隐式反馈数据非常友好。电影场景里 95% 以上的信号是播放、收藏、完播而不是显式的“用户打了多少星”。ALS 可以用置信度来处理隐式反馈用户行为越多、行为越重置信度越高。这个特性恰好和我们的用户画像标签能衔接上——我们可以把用户的画像标签转化成行为权重。第三ALS 产出的userFactors和itemFactorsuserFeature/ItemFeature天然就是用户画像里的“挖掘标签”部分能直接写回画像表形成闭环不需要单独再跑一次 Embedding 训练。3.2 从画像聚合出 ALS 训练矩阵代码逻辑与参数接下来是把画像数据转成 ALS 的训练输入。ALS 训的不是画像表而是用户-电影稀疏矩阵。这张矩阵我用的是“每个行为事件一行”的方案不做第二步压缩// ALS 训练数据构造把事件表直接映射成带权重的评分 import org.apache.spark.ml.recommendation.ALS import org.apache.spark.sql.functions._ val trainingDF eventDF .select( $uid.cast(int).as(userId), // ALS 要求 userId 是数值类型 $mid.cast(int).as(itemId), when($action play, 1.0) .when($action collect, 2.0) .when($action review, 3.0) .otherwise(0.5) .as(rating) // 隐式反馈打分权重 ) val als new ALS() .setMaxIter(10) .setRank(20) .setRegParam(0.01) .setUserCol(userId) .setItemCol(itemId) .setRatingCol(rating) .setImplicitPrefs(true) // 关键开启隐式反馈 .setAlpha(2.0) val model als.fit(trainingDF) model.write.save(/model/als_recall_v1)逻辑说明这段代码先把 UID 和 MID 转成数值 IDALS 不接受字符串 ID这是 Spark MLlib 的老限制再把行为映射成评分权重。play比collect权重低review最高因为写影评的行为损耗高于收藏。.setImplicitPrefs(true)一旦打开ALS 内部会把 rating 当成置信度而非真实评分这个开关直接决定召回偏“想要”还是偏“看过”。参数说明rank20是隐因子维度取值看数据量数据量大可以到 50数据集小 10~20 就够alpha2.0控制置信度曲线的陡峭程度隐式反馈场景建议调10~40之间试一轮maxIter10我实测跑到 15 次以后阈值变化基本在 0.001 以内不必加大。特别注意rating不要归一化到 0~1因为 ALS 的置信度机制本身会把数值当作权重你归一化反而丢失了行为强度差异。3.3 召回后过滤与融合打分如何把“画像”用起来ALS 输出的是“每个用户对一个电影的预测分”。直接按分数排序吐给用户这样不会有多少“画像”参与感而且冷启动用户会完全失效。我的三层方案是先做召回合并再做候选过滤最后做业务打分排序。召回合并阶段的思路是“多路召回归一化分数”。ALS 算一个候选集画像标签算两个候选集一是“相似电影”候选内容标签匹配召回二是“相似用户”候选把用户画像打上相似 userFactors 的用户放进来。两个候选集统一走上限 200 个候选的标准线。// 融合打分ALS 预测分 画像相似度分 import org.apache.spark.sql.functions._ val alsPred model.transform(testDF) .withColumn(als_score, $prediction) val profileScore userProfileDF .join(contentTagsDF, $profile.mid $content.mid) .groupBy($uid, $mid) .agg(avg($tag_sim).as(profile_score)) val finalRank alsPred .join(profileScore, Seq(uid, mid), left_outer) .withColumn(fusion_score, coalesce($als_score, lit(0.0)) * 0.7 coalesce($profile_score, lit(0.0)) * 0.3 ) .filter($fusion_score 0.5) .orderBy($fusion_score.desc)逻辑说明这一段做了两件事——把 ALS 预测分和画像相似度分做加权求和再把低于 0.5 的候选直接丢掉。权重 0.7/0.3 是我这个项目的设定具体取决于画像标签覆盖度如果画像覆盖率低于 60%权重会调整到 0.85/0.15。这里特别提醒ALS 的prediction在隐式反馈模式下范围通常是负数到正数和profile_score的 0~1 完全是两个量纲直接相加等于瞎搞。正确做法是先做一个 MinMax 缩放或者 Rank 变换我习惯把两路分数都转成 0~1 的排名分再融合。4. 把 Spark 调好才能出效果内存、并行度和集群部署的三个必调参数4.1 设置 executor 内存和 shuffle 分区让任务稳定跑完Spark 引入的好处是能扛住全量离线画像但一开始跑就 OOM 或疯狂 GC 的案例我见过太多次尤其是把参数交给默认值去跑数据量大的事件表时。我固定使用的参数模板是spark-submit \ --class com.example.ProfileJob \ --master yarn \ --deploy-mode cluster \ --executor-memory 8g \ --driver-memory 4g \ --executor-cores 4 \ --num-executors 8 \ --conf spark.sql.shuffle.partitions200 \ --conf spark.memory.fraction0.75 \ --conf spark.memory.storageFraction0.5 \ my-rec-1.0.jar参数说明spark.sql.shuffle.partitions200直接决定 Join 和 GroupBy 后的分区数量默认在 Spark 3.x 是 200但在数据量特别大时 200 不够用容易产生数据倾斜如果事件表在 1 亿条以上我会开到 400。spark.memory.fraction0.75是给执行内存和存储内存的共同比例剩下的 25% 留给了用户代码和元数据。storageFraction0.5控制 RDD 缓存时最多占用的那部分内存如果你的画像表要反复 Join 很多次可以把它调高到 0.6但代价是执行内存变小。踩坑经验是executor-memory不是越大越好。因为 YARN 分配容器时是按“内存 overhead”计算的8g 的 executor 实际向 YARN 申请约 10g如果你num-executors乘出来的总量超过队列上限任务会一直 Pending。4.2 ALS 训练本身的四个关键参数rank、regParam、iterations、alphaALS 是整套系统里效果最敏感的黑匣子参数波动大到我一度怀疑“换一版数据效果差一倍”。把四个必调参数讲清楚rank是隐因子维度影响的是人脸识别的“向量宽度”。用它去拟合用户口味和电影内容。20 是常见起点数据量小用 10 足够数据量大10 万用户 3 万电影可以试 50。更多维度的形式是让模型更“挑剔”但训练和预测都会变慢。regParam是正则化系数防止矩阵分解过拟合。这个参数用 0.01 起步如果训练集 MAE 低、验证集 MAE 明显偏高就把 regParam 往上涨一个数量级0.05 或 0.1。反过来验证集曲线还在明显下降就适当调低。iterations是 ALS 迭代次数。注意 Spark 社区版实现里 ALS 的maxIter次数并不是越多越好有论文支撑说 ALS 在 10 次之后基本收敛我用 15 次就明显感觉到训练时间成倍增加而中午指标没动静一般 10 次封顶。alpha是上一节提的置信度参数这里再重复讲一次它的作用是控制“行为次数与置信度的映射关系”。alpha1.0时一次观看和一百次观看的置信度差距不会太大alpha40时高频行为用户的置信度会被放大让召回更偏重深度用户。如果你的用户行为跨度很大比如既有只看过 1 部的轻度用户又有 1000 部的资深影迷我建议用 20~40 区间。4.3 数据倾斜和广播变量两个必须防的“翻车点”画像聚合最典型的翻车点是数据倾斜。电影领域“热门电影”天然集中了几个头部内容而关联这些热门电影的行为事件占了全量大头按 mid GroupBy 时某几个 reducer 处理的事件数量比其他多十倍。我处理这个问题的顺序是 第一先看 Spark UI 的 Stage 详情里最长 Task 耗时和最短 Task 耗时比是否大于 5 倍确认倾斜存在。 第二给倾斜的 key 加盐salting。把 mid 转成mid _ random(10)先打散聚合一轮再把结果按原 mid 做第二轮聚合。 第三如果只是单表很小的维度表直接广播出去。画像表经常要关联一个电影元数据表几千行这种场景下不给广播变量完全是在浪费 shuffleval broadcastMeta spark.sparkContext.broadcast(metaDF.collectAsMap()) // 之后直接在 map 函数里查映射省掉一次 shuffle join5. 避坑清单从数据清洗到交叉验证的五条血泪踩坑记录5.1 清洗后事件量少了一半问题出在 action 过滤太保守现象任务跑完画像表里的行为记录比源日志少了 50% 以上。原因在清洗阶段把action限定在非常严格的枚举里比如只留 play、collect忽视了日志里实际还有 search_click、review、follow_actor 这类行为。电影推荐场景里搜索行为往往是高强度兴趣信号漏掉它意味着大量真实兴趣被丢弃。解决行为枚举放宽至少保留 play、collect、search_click、review、share 五种。对于无法确认语义的 action统一记录到“待确认行为表”而不是直接丢弃。我当时把校验逻辑写在清洗任务里每天统计“丢弃率”并设阈值丢弃率超过 30% 则任务告警。5.2 ALS 在不同稠密度下结果差异巨大根因是隐式反馈的置信度没设现象同一套代码换了一周的数据验证集的 MAE 直接从 0.6 跳到 1.2。原因这周“只看不评”用户占比很大显式评分几乎没有ALS 在隐式反馈模式下没调alpha导致置信度完全无效。解决开启setImplicitPrefs(true)并配合alpha网格搜索。我的经验是把 alpha 在10、20、40三档上各跑一次选验证集上 Recall20 最高的一档。5.3 输出阶段 HDFS 小文件雪崩整个任务被拖垮现象画像表最终写入 HDFS 之后目录下产生了数万个小文件下一次读取时 NameNode 压力巨大任务跑一小时直接失败。原因并行度写得太高每个 task 输出一个文件实际上数据量才几 GB却被spark.sql.shuffle.partitions撑到了 200 个分区。解决最终写入前重新分区统一成 32 个分区再写profileDF.repartition(32).write.mode(overwrite).parquet(/data/profile/dt2025-01-01)写完后用fsck检查输出目录下的文件数量确保单文件大小在 128~256 MB 之间。5.4 热门电影把所有推荐结果都“带偏”了现象推荐结果里约六成都是热门头部电影个性化程度很低离线评估根本看不出来问题线上点击率一路下滑。原因ALS 公平地把热门电影推荐给了所有人因为它们在稀疏矩阵里的行item 向量稠密、置信度高、很容易被命中。解决做“热门降权”。在融合打分阶段加入一个惩罚项把计算出的“目标电影的热度分”作为乘性权重热度越高权重越低。同时要求热门电影在单用户结果中占比不超过 30%用规则硬编码切掉。5.5 冷启动用户盘点时发现没人有标签现象上线一周新注册用户比例 15%但画像命中率只有 2%。所有新用户都没法推荐。原因新用户在画像表里根本没有“统计标签”唯一的画像信息是注册信息。我把希望全押在挖掘标签上等于给没有行为的人实现了一个“也用不了”。解决冷启动阶段全部走“兴趣探测”逻辑注册页要求选三个喜欢的电影类型把选择结果直接打进画像表作为“规则标签”。这一份初始标签就能参与召回且随着行为累计统计标签逐渐接管。6. 从“能跑”到“好跑”全量物品召回、二跳反馈与画像回填的进阶技巧最后一章说三个你在基本链路跑通之后能让效果明显变好的习惯。第一把 ALS 的召回候选从“全量物品”压到“物品邻居”。全量物品做预测在百万级电影时会被躁点干扰推荐结果里低质量电影容易混进来。做法是为每部电影预计算与其隐因子最接近的 TOP50 电影用向量相似度离线算召回阶段先用用户画像锁定一组种子电影再从种子电影的邻居里挑候选。这个改造建议用一次“评估前后对比”压候选前离线 Recall20 是 0.12压完后能到 0.18 以上核心原因是候选里噪声少了模型分数在更干净的集合里排序更准。第二把二跳反馈回填到标签里。用户点了推荐列表里的电影 A然后去搜了电影 B这种行为说明系统成功把“兴趣从 A 跳到 B”是非常值得记住的信号。我在项目里会增加一个“推荐位点击”事件类型每周跑一次把这类二跳关系作为权重补进画像的相似度计算中。目的是让画像表里“相似电影”的质量不断提升而不只是依赖内容标签。第三把业务规则写成“补水层”不要写成“硬覆盖层”。比如“新人必须推热门电影”这条规则我的做法不是直接覆盖最终排序而是只在前 10 个位置里插入 2 部热门电影老用户则完全不触发热门强制推荐。业务规则的作用应该是兜底而不是越权——它一旦覆盖算法结果画像的效果就永远验证不出来。这套系统我从 T1 离线画像一路走到加上实时画像、ALS 召回、融合打分最后把画像命中率从 30% 做到 77%推荐点击率提升约 5 个百分点。最深的教训是画像和推荐绝对不能是两拨人各做一套画像资产必须直接接进召回的 signal 里否则 Spark 再快也只是给黑匣子加速。希望帮到你。本文还有配套的精品资源点击获取
返回列表