
简介这份资源是一套基于Spark的电影推荐系统完整项目包面向计算机、人工智能、通信工程等专业的在校学生与教师可用于毕业设计、课程设计、项目立项演示或自学进阶。项目整合了爬虫数据采集、Web网站、后台管理系统与Spark推荐算法四大模块并附有详细文档帮助读者理解从数据抓取到推荐落地的全流程。压缩包共1417个文件约59.61MB以html、css、js等前端页面文件java、scala、py等后端与算法源码以及png、gif、jpg等图片素材为主另含xml、json、parquet、properties等配置与数据文件目录结构清晰便于按模块检索。目前已有60人学习下载。项目源码经测试运行成功答辩评审达95分读者可据此掌握推荐系统架构设计、前后端交互与Spark算法实现并在此基础上修改扩展功能。1. 从零搭一套电影推荐系统为什么“全链路”比单点模型更值得做很多人第一次接触推荐系统都是从调一个协同过滤算法开始的拿一份评分数据跑一遍 ALS算出 RMSE觉得大功告成。但真到落地环节就会发现模型只是整条链路里最不值钱的一环——没有数据来源模型是空转的没有前端承接推荐结果没人看得到没有后台管理运营连一部电影都上不了架。这套“基于 Spark 的电影推荐系统”之所以值得认真做一遍恰恰因为它把爬虫、Web 网站、后台管理系统和 Spark 推荐引擎串成了一条完整链路是一个能写进简历、也能真正跑起来的工程闭环。它适合三类人想从“会调库”进阶到“能交付系统”的学生需要一套可演示、可扩展架构的初中级工程师以及想理解推荐系统数据流全貌的转行者。整条链路的核心逻辑是爬虫负责把电影和评分数据抓回来Spark 负责离线训练和生成推荐Web 网站负责把推荐结果呈现给用户后台管理系统负责数据维护和运营。四者缺一不可而 Spark 是把它们粘合起来的中枢。下面我按真实搭建顺序把每一步拆开讲清楚。2. 数据从哪来爬虫采集与清洗的完整落地路径推荐系统的天花板由数据决定所以第一步永远是解决“数据从哪来”。电影类数据通常包含两类电影元信息片名、类型、导演、上映年份、简介和用户行为数据评分、评论、观看记录。元信息可以靠爬虫从公开的电影资料站点采集行为数据则往往需要自己构造或从公开数据集迁移。这里我用一个通用的爬虫结构来演示具体站点结构请按目标站点的实际 DOM 调整。2.1 用 requests BeautifulSoup 抓取电影元信息先写一个最小可用的采集脚本把列表页和详情页分开处理。列表页负责拿到详情页链接详情页负责解析字段。这种“两级抓取”是爬虫里最稳的结构避免一次性把逻辑写死在一个函数里。import requests from bs4 import BeautifulSoup import time import csv import random HEADERS { User-Agent: Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0 Safari/537.36 } def fetch(url, retry3): 带重试的请求封装失败后指数退避 for i in range(retry): try: resp requests.get(url, headersHEADERS, timeout10) if resp.status_code 200: resp.encoding resp.apparent_encoding return resp.text except requests.RequestException as e: print(f第{i1}次请求失败: {e}) time.sleep(2 ** i) # 1s, 2s, 4s 退避 return None def parse_list(html): 从列表页提取详情页链接 soup BeautifulSoup(html, html.parser) links [] for a in soup.select(div.item a.title): href a.get(href) if href: links.append(href) return links def parse_detail(html): 从详情页解析电影字段字段名按目标站点调整 soup BeautifulSoup(html, html.parser) movie {} movie[title] soup.select_one(h1 span).get_text(stripTrue) if soup.select_one(h1 span) else movie[year] soup.select_one(span.year).get_text(stripTrue) if soup.select_one(span.year) else movie[rating] soup.select_one(strong.rating_num).get_text(stripTrue) if soup.select_one(strong.rating_num) else movie[genres] |.join([g.get_text(stripTrue) for g in soup.select(span.genre)]) movie[summary] soup.select_one(span[propertyv:summary]).get_text(stripTrue) if soup.select_one(span[propertyv:summary]) else return movie if __name__ __main__: base https://example-movie-site.com/top250?start{} all_movies [] for page in range(0, 250, 25): html fetch(base.format(page)) if not html: continue for link in parse_list(html): detail_html fetch(link) if detail_html: all_movies.append(parse_detail(detail_html)) time.sleep(random.uniform(1, 3)) # 随机间隔降低被封风险 print(f已采集 {len(all_movies)} 条) with open(movies.csv, w, newline, encodingutf-8) as f: writer csv.DictWriter(f, fieldnames[title, year, rating, genres, summary]) writer.writeheader() writer.writerows(all_movies)这段代码的关键点有三个。第一fetch里做了重试和指数退避网络抖动是爬虫最常见的翻车点没有重试机制跑一半就断。第二time.sleep(random.uniform(1, 3))是随机间隔固定间隔容易被识别为机器行为。第三字段解析全部用select_one加空值判断避免某个字段缺失导致整条记录报错。参数上timeout10是请求超时太小会误判失败太大则卡住流程retry3是重试次数一般 3 次足够再多说明目标站点已经不可用。2.2 数据清洗把脏数据挡在推荐引擎之外爬回来的数据几乎不可能直接进模型。常见问题包括评分字段带“分”字、年份带括号、类型字段分隔符不统一、简介里有换行和多余空格。清洗这一步做不干净后面 Spark 训练出来的推荐结果就是玄学。import pandas as pd import re def clean_movies(path): df pd.read_csv(path) # 评分去掉非数字字符转 float df[rating] df[rating].astype(str).str.extract(r(\d\.?\d*)).astype(float) # 年份提取四位数字 df[year] df[year].astype(str).str.extract(r(\d{4})).astype(Int64) # 类型统一分隔符为竖线去空格 df[genres] df[genres].astype(str).str.replace(r[\s/、], |, regexTrue).str.strip(|) # 简介去换行、去多余空格 df[summary] df[summary].astype(str).str.replace(r\s, , regexTrue).str.strip() # 丢弃关键字段缺失的行 df df.dropna(subset[title, rating]) df df[df[rating] 0] return df if __name__ __main__: df clean_movies(movies.csv) df.to_csv(movies_clean.csv, indexFalse, encodingutf-8) print(f清洗后剩余 {len(df)} 条评分范围 {df[rating].min()} - {df[rating].max()})清洗逻辑里str.extract(r(\d\.?\d*))用正则把评分里的数字抠出来比直接replace更稳因为不同站点的评分格式不一样。dropna只针对 title 和 rating因为这两个字段是推荐系统的最小依赖缺了就没法用。年份用Int64而不是int是为了保留空值而不报错。清洗完一定要打印一下数量和评分范围这是判断数据质量的第一个信号——如果评分范围是 0 到 0说明正则没匹配上得回去看原始数据长什么样。3. Spark 推荐引擎ALS 训练、调参与离线推荐的工程细节数据准备好之后就进入 Spark 的主场。电影推荐最常用的算法是 ALS交替最小二乘它属于协同过滤里的矩阵分解方法核心思想是把用户-物品评分矩阵拆成两个低维矩阵用隐向量表示用户和物品再通过内积预测缺失评分。Spark MLlib 里的ALS实现成熟、分布式、支持隐式反馈是这类系统的默认选择。3.1 用 Spark 读取数据并构建评分矩阵ALS 需要的数据格式是(userId, itemId, rating)三元组。电影元信息里的评分是“电影平均分”不是“某个用户对某电影的评分”所以真实系统里必须构造用户行为数据。常见做法有两种一是用公开的评分数据集如 MovieLens 格式迁移二是用规则模拟生成。这里演示从 CSV 读取并构建训练集的标准流程。from pyspark.sql import SparkSession from pyspark.sql.functions import col, when, rand from pyspark.ml.recommendation import ALS from pyspark.ml.evaluation import RegressionEvaluator spark SparkSession.builder \ .appName(MovieRecommender) \ .config(spark.sql.shuffle.partitions, 200) \ .config(spark.executor.memory, 4g) \ .getOrCreate() # 读取评分数据格式userId,movieId,rating,timestamp ratings spark.read.csv(ratings.csv, headerTrue, inferSchemaTrue) ratings ratings.select( col(userId).cast(int), col(movieId).cast(int), col(rating).cast(float) ).dropna() # 划分训练集和测试集 train, test ratings.randomSplit([0.8, 0.2], seed42) print(f训练集 {train.count()} 条测试集 {test.count()} 条)spark.sql.shuffle.partitions设成 200 是经验值默认的 200 在中小数据集上够用数据量再大要往上调否则 shuffle 阶段会成为瓶颈。randomSplit的seed固定住保证每次划分一致方便复现。训练集和测试集的比例 8:2 是推荐系统里比较稳的划分数据量特别小的时候可以调到 7:3。3.2 ALS 参数怎么调rank、regParam、maxIter 的取舍ALS 有三个核心参数rank隐向量维度、regParam正则化系数、maxIter迭代次数。这三个参数没有万能值但有一套可复用的调参思路。参数含义常用范围调大后的影响rank隐向量维度10 ~ 200表达能力增强但容易过拟合训练变慢regParam正则化系数0.01 ~ 1.0抑制过拟合太大导致欠拟合maxIter迭代次数10 ~ 30收敛更充分但收益递减耗时增加alpha隐式反馈置信度1.0 ~ 40.0仅隐式反馈时使用越大越信任正样本als ALS( userColuserId, itemColmovieId, ratingColrating, rank50, regParam0.1, maxIter15, nonnegativeTrue, coldStartStrategydrop ) model als.fit(train) predictions model.transform(test) evaluator RegressionEvaluator( metricNamermse, labelColrating, predictionColprediction ) rmse evaluator.evaluate(predictions) print(fRMSE {rmse:.4f})nonnegativeTrue强制隐向量非负在电影评分场景下通常能提升可解释性因为评分本身是正向的。coldStartStrategydrop是关键测试集里可能出现训练集没见过的用户或物品ALS 对这类冷启动样本会预测出 NaN不 drop 的话 RMSE 直接变成 NaN这是新手最容易踩的坑之一。RMSE 在 0.8 到 1.0 之间通常算正常低于 0.7 要警惕过拟合高于 1.2 说明模型没学到东西。3.3 生成推荐结果并落库训练完模型后要给每个用户生成 Top-N 推荐并把结果写到数据库或文件供 Web 端读取。# 为每个用户生成 Top 10 推荐 user_recs model.recommendForAllUsers(10) # 展开成 (userId, movieId, rating) 的扁平结构 from pyspark.sql.functions import explode flat_recs user_recs.select( col(userId), explode(col(recommendations)).alias(rec) ).select( col(userId), col(rec.movieId).alias(movieId), col(rec.rating).alias(score) ) # 写入 MySQL供 Web 端查询 flat_recs.write \ .format(jdbc) \ .option(url, jdbc:mysql://localhost:3306/movie_rec) \ .option(dbtable, user_recommendations) \ .option(user, root) \ .option(password, your_password) \ .option(batchsize, 1000) \ .mode(overwrite) \ .save() print(推荐结果已写入数据库)recommendForAllUsers(10)返回的是一个数组列必须用explode展开成行否则没法直接写关系型数据库。batchsize1000是 JDBC 写入的批量大小太小会导致频繁网络往返太大容易内存溢出1000 是折中值。mode(overwrite)每次全量覆盖适合离线批处理场景如果要做增量更新改成append并配合时间分区。4. Web 网站与后台管理让推荐结果真正被看见模型跑通只是完成了“后台”用户看不到推荐结果这套系统就只是个脚本。Web 网站负责展示后台管理系统负责维护两者共用同一套数据库是这套架构里最容易被低估的部分。4.1 Web 端推荐展示的最小实现Web 端通常用 Flask 或 Django 起一个轻量服务从数据库读取推荐结果渲染到页面上。核心逻辑是用户登录后拿到 userId查user_recommendations表关联movies表拿到电影详情返回给前端。from flask import Flask, render_template, session, redirect import pymysql app Flask(__name__) app.secret_key your_secret_key def get_db(): return pymysql.connect( hostlocalhost, userroot, passwordyour_password, databasemovie_rec, charsetutf8mb4, cursorclasspymysql.cursors.DictCursor ) app.route(/recommend) def recommend(): if user_id not in session: return redirect(/login) uid session[user_id] conn get_db() with conn.cursor() as cur: cur.execute( SELECT m.title, m.year, m.genres, m.summary, r.score FROM user_recommendations r JOIN movies m ON r.movieId m.movieId WHERE r.userId %s ORDER BY r.score DESC LIMIT 20 , (uid,)) movies cur.fetchall() conn.close() return render_template(recommend.html, moviesmovies)这里用JOIN把推荐结果和电影详情关联起来避免在 Python 里做二次查询。ORDER BY r.score DESC保证推荐按分数排序LIMIT 20控制单页数据量。注意session里存的是 userId不是用户名因为推荐表用的是 userId 关联用用户名会多一次查询。4.2 后台管理系统的功能边界后台管理系统不需要花哨但必须覆盖四件事电影数据的增删改查、用户管理、推荐结果查看、爬虫任务触发。用同一套 Flask 加一个管理蓝图就能实现关键是把权限控制住普通用户不能进后台。功能模块核心操作数据表电影管理新增/编辑/下架电影movies用户管理查看用户、禁用账号users推荐查看按用户查推荐列表user_recommendations任务触发手动触发爬虫和训练task_log后台的“任务触发”是最容易被忽略但最实用的功能。离线推荐不是实时的运营需要能手动跑一次爬虫或重新训练模型而不是每次都登服务器敲命令。用一个简单的任务表记录触发时间和状态配合后台按钮调用脚本即可。5. 避坑与排查这套系统最容易翻车的五个地方5.1 爬虫跑一半被封数据缺了一大块现象采集到几百条后请求全部返回 403 或空页面。原因请求频率过高或 User-Agent 单一被目标站点识别。解决加随机间隔、轮换 User-Agent、必要时用代理池并且把已采集的数据落盘支持断点续采不要每次都从头跑。5.2 ALS 预测出 NaNRMSE 直接报错现象evaluator.evaluate返回 NaN。原因测试集里有训练集没出现过的 userId 或 movieIdALS 对冷启动样本预测为 NaN。解决设置coldStartStrategydrop或者在划分数据集时保证用户和物品在训练集中都出现过。5.3 推荐结果全是同一部电影现象所有用户的 Top 10 推荐高度重合。原因数据稀疏热门电影主导了隐向量模型退化成“推荐最热”。解决提高regParam抑制热门偏置或者在评分数据里对热门电影做降权也可以引入基于内容的特征做混合推荐。5.4 Spark 任务内存溢出Executor 频繁挂掉现象训练阶段报OutOfMemoryErrorExecutor 被 kill。原因rank设得太大、数据没分区、或者spark.executor.memory不够。解决先把rank降到 50 以下检查spark.sql.shuffle.partitions是否合理必要时增加 Executor 内存并开启spark.memory.fraction调优。5.5 Web 端查询推荐结果慢到超时现象推荐页加载超过 5 秒。原因user_recommendations表没建索引或者每次查询都全表扫描。解决在userId和movieId上建联合索引推荐结果表按 userId 做分区查询时只取需要的列不要SELECT *。6. 进阶技巧把离线推荐做成可迭代的工程习惯这套系统跑通之后真正拉开差距的不是模型多复杂而是迭代效率。我自己的习惯是每次调整 ALS 参数前先把当前 RMSE、训练耗时、推荐覆盖率记到一张实验表里跑完新参数再对比。没有这张表调参就是凭感觉跑十次也说不清哪次更好。import json, time def log_experiment(params, rmse, duration, pathexperiments.jsonl): record { ts: time.strftime(%Y-%m-%d %H:%M:%S), params: params, rmse: round(rmse, 4), duration_sec: round(duration, 1) } with open(path, a, encodingutf-8) as f: f.write(json.dumps(record, ensure_asciiFalse) \n) # 用法每次训练后调用 log_experiment({rank: 50, regParam: 0.1, maxIter: 15}, rmse, 128.5)另一个实用技巧是把推荐结果做分层Top 10 走 ALSTop 10 到 Top 30 用基于类型的相似度补位。ALS 在数据稀疏时对长尾物品的预测不稳定用类型相似度兜底能明显提升推荐列表的多样性。具体做法是拿用户历史高分电影的类型分布去匹配同类型里评分靠前但 ALS 没推的电影合并去重后按分数排序。最后一个习惯每次重新训练模型前先备份当前的user_recommendations表。离线推荐是全量覆盖的一旦新模型效果变差没有备份就只能等下一次训练这个后悔药我吃过不止一次。把备份做成脚本里的一行CREATE TABLE ... AS SELECT成本极低但能救命。这套系统真正的价值不在于用了多新的算法而在于它把数据采集、模型训练、服务展示和运营管理串成了一条能持续迭代的链路。先跑通再调优最后把每次实验记录下来——这个顺序别反。希望帮到你。本文还有配套的精品资源点击获取