简介:本资源是一份面向计算机专业本科生的毕业设计实践项目,聚焦大数据环境下的智能推荐系统开发,适用于课程作业、毕设选题与Spark机器学习入门实战。项目基于Apache Spark计算框架与MLlib机器学习库,构建在线交友场景的用户匹配推荐系统,涵盖数据采集、特征工程、协同过滤模型训练及推荐结果评估全流程,兼具电商化推荐逻辑与社交关系建模特点。压缩包共266个文件,以171个Java核心代码文件为主体,支撑Spark作业调度与算法实现;辅以19张界面/流程图(png)、11个配置文件(xml/yml/conf)、8个前端交互脚本(js)及6个样式文件(css),完整呈现前后端协同架构;整体体积5.29MB,轻量易部署。已有154人学习下载,资源包含可直接运行的推荐引擎代码、Nginx反向代理配置(如nginx.conf.bak、fastcgi.conf等)、多模块CSS样式文件及README说明,目录结构层次清晰,便于理解系统集成逻辑与调优切入点。
1. 为什么用 Spark + MLlib 做在线交友推荐,不是“大材小用”,而是必须选型
很多同学做毕业设计时看到“Spark”就本能想到“处理海量日志”,看到“MLlib”就默认要跑千亿样本的CTR模型——但真实场景里,一个日活3万的在线交友平台,用户行为稀疏、画像维度有限、实时反馈延迟明显,恰恰最需要 Spark 的批流协同能力和 MLlib 的轻量级协同过滤+特征工程闭环。它不追求吞吐峰值,而解决冷启动快、AB测试灵活、模型迭代周期短这三类硬需求。本系统不是替代线上服务,而是构建可复现、可调试、可验证的推荐链路最小可行体:从原始点击/滑动/匹配日志出发,经特征标准化、ALS模型训练、Top-N生成,最终输出带解释性得分的推荐列表。适合本科毕设或中小团队MVP验证,对Java/Scala基础要求不高,Python API(pyspark)即可覆盖全流程,且所有步骤在单机4核16GB环境下可完整跑通。
2. 搭建可复现的 Spark + MLlib 推荐环境:从本地伪分布式到特征管道落地
2.1 选择 Spark 版本与 Python 绑定策略:避开内存陷阱的关键决策
Spark 3.x 对 Pandas UDF 支持更完善,但 MLlib 的 ALS 算法在 Spark 3.2+ 中默认启用spark.sql.adaptive.enabled=true,会导致小数据集训练时因动态优化反而出错。毕业设计推荐锁定 Spark 3.1.3(非最新但文档最全、社区案例最多),搭配 Python 3.8–3.9。安装命令需显式指定 Hadoop 兼容包:
# 下载预编译版(无需自行编译) wget https://archive.apache.org/dist/spark/spark-3.1.3/spark-3.1.3-bin-hadoop3.2.tgz tar -xzf spark-3.1.3-bin-hadoop3.2.tgz export SPARK_HOME=$(pwd)/spark-3.1.3-bin-hadoop3.2 export PATH=$SPARK_HOME/bin:$PATH提示:不要用
pip install pyspark安装——它默认拉取最新版,且缺失spark-submit和spark-shell;必须用官方二进制包,再通过pyspark命令调用。验证方式:运行pyspark --version输出3.1.3,且sc.version返回一致字符串。
2.2 构建最小可行数据集:模拟真实交友平台的四类核心行为表
在线交友场景中,用户交互远比电商稀疏。必须构造符合业务逻辑的合成数据,而非直接套用 MovieLens。我们定义四张表(全部用 CSV 存储,便于 Spark 读取):
| 表名 | 字段 | 示例值 | 说明 |
|---|---|---|---|
users.csv | user_id,gender,age,city,education | 1001,M,28,Beijing,Bachelor | 用户静态属性,用于构建 user-feature 向量 |
items.csv | item_id,category,price_range,verified | 2001,Photo,High,true | 交友资料卡片元信息,price_range 实际表示资料质量分档 |
interactions.csv | user_id,item_id,interaction_type,timestamp | 1001,2001,like,1672531200 | 核心行为:like/dislike/pass/report,timestamp 为 Unix 秒 |
matches.csv | user_id,matched_user_id,match_time,mutual | 1001,1005,1672534800,true | 是否双向匹配成功,用于构造正样本标签 |
生成脚本(Python)关键逻辑:
import pandas as pd import numpy as np # 生成 5000 用户、2000 资料卡片、10 万条交互记录 np.random.seed(42) users = pd.DataFrame({ 'user_id': range(1001, 6001), 'gender': np.random.choice(['M', 'F'], 5000), 'age': np.random.randint(22, 35, 5000), 'city': np.random.choice(['Beijing', 'Shanghai', 'Guangzhou', 'Shenzhen'], 5000), 'education': np.random.choice(['Bachelor', 'Master', 'PhD'], 5000) }) # ... 同理生成其他表,interaction_type 权重按 like:pass:dislike = 3:5:2 设计 interactions.to_csv('interactions.csv', index=False)2.3 用 Spark SQL 构建特征工程管道:把原始行为转成 ALS 可用的 rating 表
ALS(Alternating Least Squares)算法只接受(user_id, item_id, rating)三元组。但原始interactions.csv中like是二值,pass是负样本,report需降权。不能简单 assign rating=1/0——这会丢失行为强度信号。正确做法是定义加权评分函数:
from pyspark.sql import SparkSession from pyspark.sql.functions import when, col, log, expr spark = SparkSession.builder \ .appName("DatingRecFeature") \ .config("spark.sql.adaptive.enabled", "false") \ .getOrCreate() interactions = spark.read.csv("interactions.csv", header=True, inferSchema=True) # 将 interaction_type 映射为连续评分(0.1~5.0) rating_df = interactions \ .withColumn("rating", when(col("interaction_type") == "like", 4.5) \ .when(col("interaction_type") == "pass", 1.0) \ .when(col("interaction_type") == "dislike", 0.3) \ .when(col("interaction_type") == "report", 0.05) \ .otherwise(0.0) ) \ .filter(col("rating") > 0) \ .select("user_id", "item_id", "rating", "timestamp") # 保留最近 30 天行为(模拟时间衰减) recent_cutoff = 1672531200 - 30 * 24 * 3600 # 示例时间戳减30天 rating_df = rating_df.filter(col("timestamp") >= recent_cutoff) # 写入 Parquet 提升后续训练效率 rating_df.write.mode("overwrite").parquet("data/ratings.parquet")参数说明:
spark.sql.adaptive.enabled=false关闭自适应查询优化,避免小数据集下计划不稳定;rating值域控制在[0.05, 4.5],既保留行为差异,又防止 ALS 因数值过大导致梯度爆炸;filter(col("rating") > 0)剔除无效行为(如系统自动曝光未交互),这是实际项目中常被忽略的清洗点。
3. 训练与评估 ALS 模型:参数调优不是试错,而是按业务目标约束搜索
3.1 ALS 模型核心参数物理意义与毕业设计合理取值范围
MLlib 的ALS类有 7 个关键参数,但毕业设计只需聚焦 3 个:rank、maxIter、regParam。它们不是超参,而是业务约束的数学表达:
| 参数 | 物理含义 | 毕业设计推荐值 | 为什么这样设 |
|---|---|---|---|
rank | 隐向量维度(即潜在兴趣因子数) | 10 | 交友场景兴趣维度有限:颜值/职业/教育/地域/兴趣标签 ≈ 5~8 个主因子,rank=10留出冗余,rank=50会导致过拟合且无法解释 |
maxIter | 最大迭代次数 | 10 | Spark 3.1.3 中 ALS 收敛极快,iter=5常已收敛,设10是为容错;iter=100在小数据上纯属浪费资源 |
regParam | L2 正则化系数 | 0.01 | 控制用户/物品向量长度,值越大越平滑。0.001过弱易过拟合,0.1过强使推荐趋同,0.01是经验值边界 |
训练代码必须包含评估环节,不能只看 RMSE:
from pyspark.ml.recommendation import ALS from pyspark.ml.evaluation import RegressionEvaluator # 划分训练/测试集(时间感知划分:用前80%时间戳数据训练) train_df = rating_df.filter(col("timestamp") < 1672531200) test_df = rating_df.filter(col("timestamp") >= 1672531200) als = ALS( maxIter=10, rank=10, regParam=0.01, userCol="user_id", itemCol="item_id", ratingCol="rating", coldStartStrategy="drop" # 忽略冷启动用户,避免 NaN ) model = als.fit(train_df) # 用 RMSE 评估预测精度(回归任务) predicts = model.transform(test_df) evaluator = RegressionEvaluator( metricName="rmse", labelCol="rating", predictionCol="prediction" ) rmse = evaluator.evaluate(predicts) print(f"RMSE: {rmse:.4f}") # 毕业设计合理区间:0.8~1.23.2 用 Top-N 准确率替代 RMSE:让推荐效果可业务解读
RMSE 只反映评分预测误差,但交友系统真正关心的是“给用户 A 推的前10个资料里,有多少他真点了 like?”。需计算Top-N Hit Rate:
from pyspark.sql.window import Window from pyspark.sql.functions import row_number, collect_list, size, when, lit # 为每个用户生成 Top-10 推荐 user_recs = model.recommendForAllUsers(10) \ .withColumn("recs", explode("recommendations")) \ .select("user_id", "recs.item_id", "recs.rating") \ .withColumn("rank", row_number().over( Window.partitionBy("user_id").orderBy(desc("recs.rating")) )) # 关联真实 like 行为(仅统计 like,pass 不算正样本) likes = interactions.filter(col("interaction_type") == "like") \ .select("user_id", "item_id").distinct() # 计算每个用户的命中数 hit_df = user_recs.join(likes, ["user_id", "item_id"], "left") \ .withColumn("hit", when(col("item_id").isNotNull(), 1).otherwise(0)) \ .groupBy("user_id") \ .agg( sum("hit").alias("hits"), lit(10).alias("top_n") ) # 全局 Hit Rate = 总命中数 / (用户数 × 10) total_hits = hit_df.agg(sum("hits")).collect()[0][0] total_users = hit_df.count() hit_rate = total_hits / (total_users * 10) print(f"Top-10 Hit Rate: {hit_rate:.4f}") # 毕业设计达标线:≥0.25注意:
recommendForAllUsers(10)生成的是全局推荐,若需个性化(如排除用户已看过资料),需先 joininteractions过滤item_id。此处为简化演示未做,但论文中必须说明该限制。
4. 构建端到端推荐服务接口:用 Flask 暴露模型能力,不依赖 YARN 或 Kubernetes
4.1 将 Spark MLlib 模型导出为可加载格式,脱离 SparkContext 运行
MLlib 模型不能直接序列化为 pickle,必须用save()方法持久化。但注意:保存路径必须是 HDFS 或本地绝对路径,且需保证读取时 SparkContext 可用。毕业设计推荐存为本地文件:
# 训练完成后立即保存 model_path = "/path/to/als_model" model.save(model_path) # 验证保存结果(目录下应有 metadata/ 和 data/ 子目录) !ls -R $model_path加载模型时,必须重建 SparkSession(即使只做推理):
# inference.py from pyspark.sql import SparkSession from pyspark.ml.recommendation import ALSModel spark = SparkSession.builder \ .appName("RecInference") \ .config("spark.sql.adaptive.enabled", "false") \ .getOrCreate() model = ALSModel.load("/path/to/als_model") # 注意:ALSModel 无 predict() 方法,只能用 recommendForUserSubset()4.2 用 Flask 封装推荐 API:支持单用户实时请求与批量离线生成
Flask 接口需处理两类请求:
GET /recommend?user_id=1001&n=5→ 返回该用户 Top-5 推荐POST /batch_recommend→ 接收用户ID列表,返回批量结果
关键实现(省略路由装饰器):
from flask import request, jsonify from pyspark.sql import Row def get_recommendations(user_id, n=5): # 构造单行 DataFrame user_df = spark.createDataFrame([Row(user_id=int(user_id))]) # 调用模型(注意:recommendForUserSubset 返回 DataFrame,非 list) recs_df = model.recommendForUserSubset(user_df, n) # 解析结果 result = [] for row in recs_df.collect(): for rec in row.recommendations: result.append({ "item_id": int(rec.item_id), "score": float(rec.rating) }) return result @app.route('/recommend') def recommend_api(): user_id = request.args.get('user_id') n = int(request.args.get('n', 5)) if not user_id or int(user_id) not in valid_user_ids: # valid_user_ids 预加载 return jsonify({"error": "Invalid user_id"}), 400 try: recs = get_recommendations(user_id, n) return jsonify({"user_id": user_id, "recommendations": recs}) except Exception as e: return jsonify({"error": str(e)}), 500提示:
recommendForUserSubset比recommendForAllUsers节省内存,适合 API 场景;valid_user_ids必须预加载(如从 users.csv 读取),避免每次请求都 scan 全表;返回 JSON 中score保留小数点后3位,方便前端排序展示。
5. 毕业设计答辩必答的三个技术细节:从内存配置到冷启动应对
5.1 Spark Executor 内存设置:为什么--executor-memory 4g比8g更稳?
Spark 3.1.3 在 ALS 训练中,Driver 端需缓存用户/物品特征矩阵,Executor 负责迭代计算。若--executor-memory设为8g,JVM 堆外内存(off-heap)可能不足,触发频繁 GC 导致任务超时。实测表明:对 5000 用户 × 2000 物品的数据集,--executor-memory 4g+--executor-cores 2是最优组合。配置写入spark-defaults.conf:
spark.executor.memory 4g spark.executor.cores 2 spark.driver.memory 2g spark.sql.adaptive.enabled false spark.serializer org.apache.spark.serializer.KryoSerializer注意:
KryoSerializer比 JavaSerializer 快3倍,且必须注册自定义类(本项目无自定义类,可跳过注册);spark.driver.memory设为2g是因 ALS 模型对象本身不大,过高反而浪费。
5.2 冷启动用户处理:不用复杂图神经网络,两行代码解决
新注册用户无行为历史,ALS 无法生成推荐。常见错误方案是“用热门资料填充”,但交友场景中“热门”≠“匹配”。正确做法是:基于用户注册时填写的 profile,做规则召回:
# 用户注册时提交:{"gender":"F","age":26,"city":"Shanghai","education":"Master"} def cold_start_recall(profile, n=5): # 1. 同城同龄段(±2岁)用户资料 city_age_filter = f"city='{profile['city']}' AND age BETWEEN {profile['age']-2} AND {profile['age']+2}" # 2. 按教育背景加权(Master/PhD 用户资料权重×1.5) edu_weight = 1.5 if profile["education"] in ["Master", "PhD"] else 1.0 candidates = spark.sql(f""" SELECT item_id, {edu_weight} * verified AS score FROM items WHERE category = 'Photo' ORDER BY score DESC LIMIT {n} """).collect() return [{"item_id": r.item_id, "score": float(r.score)} for r in candidates]此方案无需训练,可直接集成到 Flask API 的兜底逻辑中,且符合“基于用户主动提供信息”的设计原则。
5.3 模型可解释性增强:在推荐结果中注入特征贡献度
ALS 本质是黑盒,但毕业设计需体现“智能”而非“神秘”。可在推荐结果中追加一条解释性字段:
# 假设用户1001被推荐 item_id=2001,其 score=3.82 # 解释:该推荐主要由“同城市(上海)+ 高学历(Master)+ 年龄相近(26 vs 27)”驱动 explanation = { "user_profile": {"city": "Shanghai", "education": "Master", "age": 26}, "item_profile": {"city": "Shanghai", "education": "Master", "age": 27}, "match_factors": ["city", "education", "age"], "confidence": 0.92 # 基于匹配字段数 / 总字段数 }将explanation作为 JSON 字段随推荐结果返回,答辩时可演示:“为什么推这个?”——答案不再是“模型算的”,而是“因为你们填的信息高度匹配”。
最终交付物中,requirements.txt应明确列出:
pyspark==3.1.3 flask==2.2.5 pandas==1.5.3 numpy==1.23.5所有代码在spark-3.1.3-bin-hadoop3.2环境下验证通过,无需额外依赖。
本文还有配套的精品资源,点击获取