资讯动态

基于Spark的菜品推荐系统:从评分预处理到ALS协同过滤实战

发布时间:2026/9/11 14:31:28 来源:尧图企业网站定制
简介基于Spark的餐饮平台菜品智能分析推荐系统源码与数据库源自个人毕业设计答辩评审达98分。项目围绕菜品数据采集、分析与推荐展开覆盖Spark处理流程、推荐算法落地及前端展示适合计算机、通信、人工智能、自动化等相关专业学生用于课程设计或毕业设计也可供入门及进阶开发者研习。压缩包共49个文件以17个Java源码、8个XML配置、6个CSS、5个JS为主另含JSP页面、SQL数据库脚本、CSV与JSON数据样例及README说明文档整体约2.05MB结构清晰便于按模块查阅。已有254人学习下载是理解餐饮推荐系统完整实现路径的实用参考资料。1. 评分数据过万之后菜品推荐为什么必须交给 Spark拿到这份餐饮评分工程时我第一反应不是去看推荐算法本身而是先翻数据目录里那份user_meal_rating.csv。原因很简单当评分数据量过万、又带着时间戳和菜品维度时单机 pandas 的 join 和 groupBy 会越来越吃力内存溢出往往出现在你以为最稳的聚合步骤里。这个项目把「数据清洗 → 协同过滤 → Java 接口」串成一条完整链路核心计算全部落在 Spark 上对要交毕设或者想补 Spark 实战的人来说参考价值基本等同于一堂完整的分布式推荐系统课。它可以解决三类问题怎么用 Spark SQL 和 DataFrame 做评分数据的特征工程怎么用 MLlib 的 ALS 算法产出菜品推荐以及 Spark 计算结果如何回写 MySQL 并由 Java 服务对外暴露接口。适合计算机、大数据相关专业的学生也适合刚接触 Spark、想看到一个真实项目里各组件如何咬合的一线工程师。2. 数据模型与预处理把 user_meal_rating 变成可训练特征很多人在推荐系统上翻车不是算法选错而是原始数据没有经过合理的建模。user_meal_rating这类评分表看似简单实际涉及字段类型推断、空值分布、菜品长尾效应三个问题任何一个处理不当ALS 训练出来的 latent factor 都会失真。这一章我会按“建表 → 探查 → 过滤”的顺序把预处理阶段的关键代码拆开讲。2.1 Spark SQL 建表与 CSV/JSON 双格式读取工程里同时提供了user_meal_rating.csv和user_meal_rating.json两份数据这在高分项目里很常见CSV 方便直接导入 MySQLJSON 方便 Spark 原生读取。先用 Spark SQL 把 CSV 注册成表后续所有探查和特征工程都能用 SQL 表达。CREATE DATABASE IF NOT EXISTS meal_recommend; USE meal_recommend; CREATE TABLE IF NOT EXISTS raw_rating ( user_id INT, meal_id INT, rating DOUBLE, rating_time STRING ) USING CSV OPTIONS ( path /data/user_meal_rating.csv, header true, inferSchema true, sep , );这段建表语句的关键在USING CSV和OPTIONS。USING CSV是 Spark 2.0 之后推荐的外部数据源写法数据文件可以放在 HDFS、S3 或本地路径Spark 执行查询时才真正读取文件内容而不是像传统数据库那样先导入再计算。header true告诉解析器跳过首行字段名inferSchema true让 Spark 自动推测每列类型。如果不开 inferSchemarating会被当成字符串后面做avg(rating)时要么报错要么默默丢数据。如果偏好 DataFrame APIJSON 版本可以这样读from pyspark.sql import SparkSession spark SparkSession.builder \ .appName(meal-preprocess) \ .master(local[*]) \ .getOrCreate() rating_df spark.read \ .option(inferSchema, true) \ .json(/data/user_meal_rating.json) rating_df.printSchema() rating_df.show(5, truncateFalse)printSchema()输出的结果里user_id和meal_id会被推成 integerrating推成 double。看到rating_time是字符串而不是 timestamp 时不要急着转换原始时间字段经常带时区或者前后缀先保留字符串做探查更安全。这个阶段的核心目标只有一个让 DataFrame 的 schema 和后续 ALS 要求的userCol、itemCol、ratingCol对齐。2.2 评分分布探查与倾斜数据过滤数据表建好之后先别急着训练。我习惯先用describe和groupBy探查评分分布这一步能直接决定要不要过滤冷门菜品、要不要做数据分层。2.2.1 按菜品聚合统计热度from pyspark.sql import functions as F rating_df.describe([rating]).show() meal_stats rating_df.groupBy(meal_id).agg( F.count(user_id).alias(cnt), F.avg(rating).alias(avg_rating), F.stddev(rating).alias(std_rating) ) meal_stats.orderBy(F.desc(cnt)).show(10)describe([rating])输出的是评分列的 count、mean、stddev、min、max重点看 count 和原表行数是否一致差异部分就是空值。groupBy(meal_id)这段是标准的 spark 数据分析案例写法cnt是每个菜品的评分次数avg_rating是平均分std_rating能看出评分分歧度。真实项目里最容易忽略的是菜品热度分布。餐饮数据的评分高度向头部菜品集中少数爆款菜品占据大部分评分大量小众菜品只有三五条记录。这种数据直接丢给 ALS评分次数过少的菜品会被隐因子压缩成噪声最终推荐列表里反复出现热门菜用户看到的「个性化」实际上是「平均化」。2.2.2 过滤低评分次数菜品popular_meals meal_stats.filter(F.col(cnt) 5).select(meal_id) train_base rating_df.join(popular_meals, onmeal_id, howinner) print(fbefore: {rating_df.count()}) print(fafter: {train_base.count()})过滤阈值我一般取 5也就是菜品至少有 5 个用户评过分才进入训练集。阈值太小起不到过滤作用阈值太大比如 20会丢掉大量长尾菜品让推荐结果失去多样性。join用inner模式实现的是“只在保留菜品列表内取评分记录”的语义。运行后对比 before 和 after 的行数差就能估算数据倾斜程度。这里有个容易踩的坑过滤操作要在切分训练集和测试集之前做否则测试集里会出现训练阶段从未见过的 meal_id。ALS 对未知物品做预测时会输出空值虽然可以靠coldStartStrategydrop兜底但测试集样本被丢掉会直接影响 RMSE 的可信度。preprocess 做完后把train_base缓存起来下一步就是特征工程与模型训练。3. ALS 协同过滤菜品推荐的召回核心预处理完成后进入推荐核心。Spark 的 MLlib 里最适合这种评分场景的算法是 ALS交替最小二乘它属于协同过滤里的矩阵分解家族。菜品推荐用协同比基于内容推荐更合理的原因在于菜品的口味、辣度、甜度这些属性很难用结构化标签穷举但用户的历史评分行为天然隐含了这些偏好。3.1 为什么选 ALS 而不是手工统计推荐手工统计推荐的做法是算每个菜品的平均分然后排序输出实现简单但没有个性化。ALS 的思想是把用户 × 菜品的评分矩阵拆成两个低秩矩阵U和V让U_i · V_j尽量逼近rating_ij。拆出来的隐因子虽然没有显式语义但实际训练后往往能对应到“口味浓郁程度”“是否偏辣”“价格敏感度”等潜在维度。ALS 的另一个优势是天然支持分布式交替固定U更新V、固定V更新U每一步都是独立的矩阵运算可以并行跑在多个 executor 上。3.2 参数怎么设、RMSE 怎么看ALS 需要调节的核心参数不多但每个都直接影响结果质量。这个项目里我通常先跑一组基线参数再根据验证集的 RMSE 微调。参数作用建议范围rank隐因子个数决定模型表达能力8 ~ 20regParam正则化系数防止过拟合0.01 ~ 0.2maxIter最大迭代轮数收敛控制5 ~ 15alpha隐式反馈置信度权重仅隐式场景10 ~ 40coldStartStrategy冷启动处理drop / nanrank 过小模型欠拟合菜品差异学不出来rank 过大训练耗时长且容易过拟合到评分噪声。regParam的作用是抑制隐因子向量的模长评分数据稀疏时适当调大。maxIter不是越大越好我见过很多项目迭代到 20 轮以后 RMSE 几乎不动白白浪费 spark 内存和计算时间。3.2.1 训练代码与验证from pyspark.ml.recommendation import ALS from pyspark.ml.evaluation import RegressionEvaluator (train, test) train_base.randomSplit([0.8, 0.2], seed42) als ALS( userColuser_id, itemColmeal_id, ratingColrating, rank12, regParam0.08, maxIter10, coldStartStrategydrop ) model als.fit(train) evaluator RegressionEvaluator( metricNamermse, labelColrating, predictionColprediction ) rmse evaluator.evaluate(model.transform(test)) print(frmse{rmse:.4f})randomSplit按 8:2 切分训练集和测试集seed42是为了结果可复现。coldStartStrategydrop表示测试集里遇到训练时没见过的 user 或 meal 时直接丢弃预测结果而不是输出 NaN。RMSE 的合理区间和评分尺度强相关rating 范围是 1 到 5 分时RMSE 在 0.9 以下基本可用低于 0.7 说明模型已经能把用户偏好学得比较清楚。如果基线 RMSE 降不下来优先调rank而不是maxIter。常见做法是写一层循环让 rank 取 8、12、16各跑一次并记录 RMSE选最低的那个。整个调参过程在 Spark 上跑很快因为 ALS 的迭代逻辑在 MLlib 内部是并行化的。3.3 生成并展开 Top-N 推荐结果训练完成的model可以直接为所有用户生成推荐列表这一步产出的是召回结果后续 Java 服务层只需要按 user_id 查询。user_recs model.recommendForAllUsers(10) user_recs.select(user_id, recommendations).show(5, truncateFalse) from pyspark.sql import functions as F rec_exploded user_recs.select( user_id, F.explode(recommendations).alias(rec) ).select( user_id, F.col(rec.meal_id).alias(meal_id), F.col(rec.rating).alias(score) ) rec_exploded.orderBy(user_id, F.desc(score)).show(20)recommendForAllUsers(10)返回的结构是user_id recommendations 数组数组里每个元素是(meal_id, rating)结构体。explode把数组拆成多行然后取出meal_id和预测评分score。这一步得到的是带分数的完整推荐明细可以直接回写到 MySQL也可以先落成 Parquet 供后续离线分析。需要提醒的是recommendForAllUsers产出的 score 不是真实评分而是模型预测的偏好强度。两个菜品的 score 差 0.1 不代表用户真的会给出 0.1 分的差异所以在 Java 层展示时我一般只取排序关系不直接展示预测分。4. Java 服务层接入Spark 结果落库与 API 暴露Spark 算完只是第一步用户能感知到的是最终接口返回的菜品列表。这个项目的 Java 部分承担两个职责把 Spark 离线产出的推荐结果写入 MySQL以及提供一个 HTTP 接口按用户 ID 查询推荐菜品。这也是「智能分析推荐系统」里「系统」二字的落点。4.1 spark.sql 结果表设计与数据回写工程根目录里的spark.sql包含了建表语句实际设计时我建议把推荐结果表和菜品热度统计表分开避免单表数据量膨胀后影响查询效率。CREATE TABLE IF NOT EXISTS result_recommend ( id INT AUTO_INCREMENT PRIMARY KEY, user_id INT NOT NULL, meal_id INT NOT NULL, score DOUBLE, rank_no INT, UNIQUE KEY uk_user_meal (user_id, meal_id) ); INSERT INTO result_recommend (user_id, meal_id, score, rank_no) VALUES (?, ?, ?, ?) ON DUPLICATE KEY UPDATE score VALUES(score), rank_no VALUES(rank_no);result_recommend每行代表「某个用户对某个菜品的预测偏好」rank_no是当前用户下的排序位次方便接口层直接按位次取数。ON DUPLICATE KEY UPDATE是用来处理 Spark 作业重跑的同一对user_id meal_id再次写入时直接更新分数而不是插入重复行。回写 MySQL 的常见做法是将上一步rec_exploded的结果写成 JDBC 批量插入。用 PySpark 可以这样实现rec_exploded.write \ .mode(overwrite) \ .jdbc( urljdbc:mysql://localhost:3306/meal_recommend?useSSLfalseserverTimezoneAsia/Shanghai, tableresult_recommend, properties{ user: root, password: your_password, driver: com.mysql.cj.jdbc.Driver } )mode(overwrite)表示作业重跑时先清空再写入但这会破坏ON DUPLICATE KEY UPDATE的语义我一般改成mode(append)配合结果表里的唯一键完成更新。JDBC 写入的批次大小受spark.sql.shuffle.partitions影响默认 200 个分区同时写 MySQL 会把数据库连接池打满建议回写前用coalesce(8)收敛分区数。4.2 Spring Boot 推荐接口与冷启动兜底回写完成后Java 侧只需要一个 Controller 和一条查询 SQL。下面是典型的 Spring Boot 实现片段。RestController RequestMapping(/api/meal) public class MealRecommendController { private final JdbcTemplate jdbcTemplate; public MealRecommendController(JdbcTemplate jdbcTemplate) { this.jdbcTemplate jdbcTemplate; } GetMapping(/recommend/{userId}) public ListMapString, Object recommend(PathVariable Integer userId, RequestParam(defaultValue 10) Integer topN) { ListMapString, Object hits jdbcTemplate.queryForList( SELECT meal_id, score FROM result_recommend WHERE user_id ? ORDER BY score DESC LIMIT ?, userId, topN ); if (hits.isEmpty()) { return jdbcTemplate.queryForList( SELECT meal_id, avg_rating AS score FROM dish_stats ORDER BY avg_rating DESC LIMIT ?, topN ); } return hits; } }JdbcTemplate的queryForList返回结构化的 Map 列表Spring Boot 会自动序列化成 JSON。第一条 SQL 从result_recommend里查指定用户的 Top-N 推荐ORDER BY score DESC保证取到的是偏好强度最高的菜品。接口层不直接接 Spark而是查 MySQL好处是响应时间从秒级降到毫秒级避免每次请求都拉起一个 SparkSession。4.2.1 冷启动用户兜底逻辑hits.isEmpty()分支就是冷启动处理。新注册用户没有历史评分result_recommend里自然没有他的推荐记录。这时候返回热门菜品列表是最稳妥的 fallbackdish_stats表由预处理阶段的meal_stats结果落库得到相当于全局平均分排行。兜底逻辑虽然简单但能保证接口永远不会返回空数组这是推荐系统上线的基本要求。评分数据积累后再跑一次 Spark 批处理用户的新偏好就会被并入下一轮推荐结果。5. 部署调优spark-submit 参数与三个常见坑源码在本地能跑通和集群上能稳定运行是两回事。把 Spark 作业提交到 YARN 集群时资源参数不合理会导致任务卡在运行中、executor 被 kill、或者 ALS 迭代十几轮后直接 OOM。我平时做 spark 集群搭建后的第一个验证任务就是用这套菜品推荐工程做基准测试。5.1 集群提交参数怎么看spark-submit \ --master yarn \ --deploy-mode cluster \ --driver-memory 2g \ --driver-cores 1 \ --executor-memory 4g \ --executor-cores 2 \ --num-executors 4 \ --conf spark.default.parallelism8 \ --class com.example.MealBatchJob \ meal-recommender-1.0.jarnum-executors是 executor 个数executor-memory是每个 executor 的堆内内存executor-cores是每个 executor 占用的 CPU 核数。这个配置给 4 个 executor、每个 4g 内存和 2 核适合百万级评分数据数据量更大时优先加num-executors而不是单 executor 内存因为 JVM 堆超过 8g 后 GC 停顿会明显拖慢 ALS 迭代。spark.default.parallelism控制的是 shuffle 后默认分区数。ALS 在每轮迭代里都会做矩阵内积和梯度聚合分区数太小时 shuffle 数据量集中到少数 task 上容易造成内存倾斜分区数太大时调度开销反而盖过计算收益。经验值是 executor 总核数的 2 到 3 倍这里 4 个 executor 乘 2 核等于 8所以配了 8。5.2 三个典型的排错点现象常见原因处理建议ALS 迭代不收敛RMSE 一直高数据倾斜严重头部菜品评分过多过滤cnt 5的菜品对meal_id做加盐分桶推荐结果全是热门菜个性化不足冷启动用户走了兜底逻辑rank 太小兜底结果按user_id缓存个性化调大 rankSpark 作业显示成功但 MySQL 无数据JDBC 写入模式错误或分区过多检查是否为append写入前coalesce(8)数据倾斜是最容易忽略的一个。餐饮场景下爆款菜品的评分可能是小众菜品的几百倍ALS 训练时这些高热度行会集中在少数 partition 里表现为部分 executor 内存暴涨、其他 executor 闲置。过滤低评分菜品能缓解一部分更彻底的做法是训练前对meal_id做哈希加盐拆分成多个子键参与计算训练完再合并回原始键。另一个实战细节是提交作业时用--conf spark.driver.maxResultSize2g限制 driver 拉取的数据量否则recommendForAllUsers结果集过大时driver 会 OOM。这类问题多发生在数据量超过千万级之后调试时把并行度和 driver 内存同时调上去比单方面加大 executor 内存更有效。本文还有配套的精品资源点击获取

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价