十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

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

基于Spark的餐饮菜品推荐:ALS协同过滤实践

基于Spark的餐饮菜品推荐:ALS协同过滤实践 简介基于Spark的餐饮平台菜品智能分析推荐系统源码数据库是一份面向高校学生的毕业设计级项目答辩评审分高达98分代码经调试可正常运行。它主要服务于计算机、大数据、人工智能、自动化等相关专业的学生与从业者适合作为期末课程设计、课程大作业或毕业设计的参考也可在此基础上二次改造实现个性化功能。包体共49个文件压缩包仅2.05MB。其中17个Java源文件构成系统核心逻辑8个XML与3个JSP负责配置和页面交互CSS/JS用于前端展示SQL、CSV及JSON文件提供数据库结构与用户菜品评分数据另有字体图标等静态资源支撑界面效果。这套源码完整覆盖从数据导入、Spark处理分析到结果推荐的流程并附带数据库脚本和样例数据便于直接运行和复现。已有254人学习下载适合希望快速理解餐饮推荐系统整体架构、掌握Spark实际应用的中级学习者。1. 基于 Spark 的菜品推荐到底解决什么问题一家餐饮平台每天产生几十万条订单、上千道菜品运营最常提的需求是“每个用户打开菜单前三个菜就是他愿意点的”。全局热销榜解决不了这件事爱吃辣的工程师和正在减脂的上班族打开菜单后应当看到完全不同的第一屏。把用户行为转成偏好、用算法生成个性化列表、再落回数据库供接口读取这是基于 Spark 的餐饮平台菜品智能分析推荐系统真正在做的事也是这类项目被当作“高分项目”的原因——它同时覆盖了数据处理、算法建模、数据库设计三块硬功夫。这套方案适合两类读者正在做数据库课程设计或推荐系统项目的学生需要一套讲得清原理、跑得通全链路的实现以及业务中要搭第一版离线推荐的工程师想了解 Spark 集群上 ALS 协同过滤的可靠落地路径。Spark 在这里的价值不是算法多新奇而是把订单日志清洗、模型训练、批量推理放进同一套分布式框架。下文按数据模型、ALS 训练、数据库回写、效果验证四条线展开。2. 菜品推荐的数据模型把订单日志翻译成评分矩阵整套系统从数据库开始最后也回到数据库。ALS 模型不关心菜名、价格和分类它只认(user_id, dish_id, rating)三元组。所以在写任何 Spark 代码之前先把数据流画清楚源表有哪些、清洗后落到哪张表、评分怎么构造。这一步想不明白后面调参再久都救不回来。2.1 餐饮订单数据的三个事实首先要接受三个现实它们决定了建模方式。第一用户没有给菜打分的习惯。美团、饿了么上真正写评价的是极少数订单直接拿dish_ratings表喂模型样本量会小到无法训练。第二订单数据量大但用户对菜品的覆盖非常稀疏。用户常吃的菜可能只有十几道平台却有上千道菜品这意味着评分矩阵里绝大多数位置是空的需要靠交替最小二乘这类矩阵分解方法补全。第三餐饮有天然复购行为同一道菜会被同一个用户反复下单这与电影、图书的一次性消费不同处理时要考虑重复购买带来的偏好强度而不是把每次下单当成独立事件。这三个事实指向同一个结论纯靠显式评分不可行订单行为才是最重要的数据资产。2.2 两种评分构造方式与适用场景维度显式评分隐式反馈数据来源dish_ratings打分表ordersorder_items行为日志评分含义用户主观喜好1~5 分购买频次、金额、最近消费时间折算样本稠密程度极稀疏绝大多数菜品无评分相对稠密下单行为天然产生对应算法implicitPrefsFalseimplicitPrefsTrue餐饮场景适配适合已做会员评价体系的平台更适合只有订单数据的普通餐饮系统我一般会优先做隐式反馈。原因是行为数据不需要额外的采集成本而且隐式 ALS 对“没买过”可以处理成 0 而不是缺失值能利用负样本信息。显式评分更适合作为项目答辩时的对比实验展示两种模式下的指标差异。2.3 用 Spark DataFrame 从订单表计算评分这里给出最核心的清洗代码。输入是 MySQL 里的orders和order_items两张表输出是 Spark DataFrame 格式的评分表列名为user_id、dish_id、rating完全对齐 ALS 的输入要求。from pyspark.sql import SparkSession from pyspark.sql.functions import sum, log1p, round spark SparkSession.builder \ .appName(restaurant-dish-rating) \ .master(yarn) \ .config(spark.sql.shuffle.partitions, 200) \ .getOrCreate() jdbc_props { user: spark_user, password: ******, driver: com.mysql.cj.jdbc.Driver } orders spark.read.jdbc( urljdbc:mysql://192.168.10.11:3306/restaurant, tableorders, propertiesjdbc_props ) order_items spark.read.jdbc( urljdbc:mysql://192.168.10.11:3306/restaurant, tableorder_items, propertiesjdbc_props ) # 只取近90天订单保留行为的新鲜度 ratings order_items.join(orders, order_id) \ .filter(order_time 2024-01-01) \ .groupBy(user_id, dish_id) \ .agg( sum(quantity).alias(buy_cnt), sum(amount).alias(total_amount) ) \ .withColumn(rating, round( log1p(buy_cnt) 0.1 * log1p(total_amount), 4 )) \ .select(user_id, dish_id, rating)这段代码的核心是最后的评分折算公式log1p(buy_cnt) 0.1 * log1p(total_amount)。用log1p是为了压缩长尾买 1 次和买 50 次的差异不是 50 倍而是被压缩到约 3.9 的差距金额前面乘 0.1是把“客单价高”的影响压制在购买次数之下防止贵价菜凭金额霸榜。filter里的时间窗口决定了推荐的新鲜度做离线评估时一定要保证训练数据的时间全部早于测试数据否则会引入前视偏差。2.4 为什么不用 RDD 自己写矩阵分解很多教程还在用 RDD 版本的ALS.train但那套 API 在 Spark 2.x 之后进入维护模式。现在的标准做法是pyspark.ml.recommendation.ALS基于 DataFrame 封装配合 Pipeline 可以一键完成特征列选择、模型保存、热加载。spark之dataframe这个关键字背后的核心差异在于DataFrame 有 Catalyst 优化器join 和 groupBy 会自动选择 broadcast join 或 sort merge join这在处理百万级订单时比手写 RDD map 快一个量级。评分表构造好之后下一步就是把它喂给 ALS 训练器。3. 用 Spark ALS 训练菜品推荐模型完整可复现代码菜品推荐的数据量级决定了 ALS 是默认选择。Item-based CF 在菜品数量上万后相似度矩阵的存储和更新都很笨重SVD 对稀疏矩阵的收敛速度不如 ALS。MLlib 内置的 ALS 把最小二乘求解的并行化封装好了训练、预测、保存模型都是几行代码的事。3.1 ALS 的迭代逻辑与隐式反馈的差异ALS 的核心是把用户-菜品评分矩阵R(m×n)近似分解成用户隐因子矩阵U(m×k)和菜品隐因子矩阵V(n×k)的乘积。直接解这个分解是非凸优化但固定住V时求解U变成了一组独立的最小二乘问题可以并行固定住U再解V同理。两个步骤交替进行直到损失收敛。餐饮场景用隐式反馈时损失函数和显式模式有本质区别显式 ALS 只对观测到的评分做误差回传隐式 ALS 会对所有零值也计算置信度权重1 alpha * rating。这意味着没买过的菜也参与了模型学习alpha 越大“高评分行为”的置信度越高。爱吃辣的用户点过剁椒鱼头 20 次模型会把香辣鸡杂排在酸菜鱼前面即使他从来没点过这道菜。3.2 训练与基础评估代码from pyspark.sql import SparkSession from pyspark.ml.recommendation import ALS from pyspark.ml.evaluation import RegressionEvaluator from pyspark.sql.functions import explode, col spark SparkSession.builder \ .appName(dish-recsys-als) \ .master(yarn) \ .config(spark.sql.shuffle.partitions, 200) \ .getOrCreate() ratings spark.read.jdbc( urljdbc:mysql://192.168.10.11:3306/restaurant, tableratings, properties{user: spark_user, password: ******, driver: com.mysql.cj.jdbc.Driver} ) # 按用户做分层切分避免同一用户数据同时落在训练和测试 train, test ratings.randomSplit([0.8, 0.2], seed42) als ALS( userColuser_id, itemColdish_id, ratingColrating, implicitPrefsTrue, rank20, maxIter10, regParam0.08, alpha40, coldStartStrategydrop ) model als.fit(train) # 给测试集中的用户-菜品对打分 pred model.transform(test)randomSplit是最简单的切分方式但把它改成按user_id做哈希分桶更严谨可以防止同一用户的订单被切到两侧造成评估乐观。coldStartStrategydrop表示测试集里出现训练集从未见过的用户或菜品时直接丢弃该行而不是输出 NaN否则下游评估指标会全部变成 NaN。3.3 五个必调参数参数默认值餐饮场景建议说明rank1020~30隐因子维度。菜品覆盖越杂、用户群越大取值越高但过大容易过拟合maxIter1010~15交替迭代次数。餐饮数据量级下 10 次足够收敛继续加只能看到损失在万分位波动regParam0.10.05~0.1L2 正则系数。评分越稀疏正则要越大防止冷门菜隐因子被噪声带偏alpha1.030~40隐式反馈置信度斜率。表示购买次数每增加一次置信度增幅的陡峭程度implicitPrefsfalsetrue与 2.3 节的评分构造方式严格对应改错会直接改变损失函数语义这里最容易踩的坑是评分表用订单行为构造却忘了把implicitPrefs设为 true。两者一旦错位模型会把“购买次数为 0”当成真实缺失值而不是隐式负样本推荐质量会明显下降。3.4 生成 Top-N 推荐并拍平输出user_recs model.recommendForAllUsers(10) results user_recs.select( col(user_id), explode(col(recommendations)).alias(rec) ).select( col(user_id), col(rec.dish_id).alias(dish_id), col(rec.rating).alias(score) )recommendForAllUsers(10)返回的结构是每行一个用户recommendations列里是一个包含 dish_id 和预测分值的结构体数组。explode会把数组拍平成行一行一个推荐菜品这样写进数据库时直接就是规整的三列。注意这个操作会产生全量用户的笛卡尔展开数据量是用户数乘以 10对 spark 内存的占用集中在 shuffle 阶段建议把spark.sql.shuffle.partitions设置成集群并发度的两到三倍比如 3 个执行核心各跑 64 个并发就设 200 左右。4. 数据库表设计与推荐结果回写的工程细节推荐模型产出的结果如果不落库就无法支撑线上接口查询。这一章解决数据库设计的问题哪些表承担源数据、哪些表承担结果存储、Spark 写回时用什么模式最安全。4.1 六张核心表的 Schema 设计CREATE TABLE users ( user_id INT PRIMARY KEY AUTO_INCREMENT, user_name VARCHAR(64) NOT NULL, city VARCHAR(32) DEFAULT NULL, register_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP, KEY idx_city (city) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4; CREATE TABLE dishes ( dish_id INT PRIMARY KEY AUTO_INCREMENT, dish_name VARCHAR(128) NOT NULL, category_id INT NOT NULL, price DECIMAL(8,2) NOT NULL, status TINYINT NOT NULL DEFAULT 1, KEY idx_category (category_id) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4; CREATE TABLE orders ( order_id BIGINT PRIMARY KEY, user_id INT NOT NULL, store_id INT NOT NULL, order_time TIMESTAMP NOT NULL, total_amount DECIMAL(10,2) NOT NULL, KEY idx_user_time (user_id, order_time) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4; CREATE TABLE order_items ( id BIGINT PRIMARY KEY AUTO_INCREMENT, order_id BIGINT NOT NULL, dish_id INT NOT NULL, quantity INT NOT NULL DEFAULT 1, amount DECIMAL(8,2) NOT NULL, KEY idx_order (order_id), KEY idx_dish (dish_id) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4; CREATE TABLE dish_ratings ( user_id INT NOT NULL, dish_id INT NOT NULL, rating FLOAT NOT NULL, rating_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP, PRIMARY KEY (user_id, dish_id) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4; CREATE TABLE recommendations ( user_id INT NOT NULL, dish_id INT NOT NULL, score DOUBLE NOT NULL, rank_no INT NOT NULL, model_version VARCHAR(32) DEFAULT , update_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, PRIMARY KEY (user_id, dish_id), KEY idx_user_rank (user_id, rank_no) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;recommendations表是这个项目的关键设计主键是(user_id, dish_id)重算时可以用 INSERT ON DUPLICATE KEY UPDATE 语义覆盖rank_no存菜品在用户列表里的序号线上接口直接WHERE user_id? AND rank_no10 ORDER BY rank_no就能取数model_version记录本次结果是哪个参数版本产出的这在对比实验时非常有用。4.2 为什么推荐结果要单独建表而不是实时计算ALS 训练出的U和V矩阵存在 Spark 里但线上应用不可能每来一个请求就拉起一个 Spark 任务。常见做法是每天凌晨用 crontab 或 DolphinScheduler 触发一次离线重算把 Top-N 结果批量写进 MySQL白天应用层只做普通 SQL 查询。这个架构下数据库课程设计里讲的“表关系设计、索引优化”就有了真实落点recommendations表只有几万到几十万行全表扫描都能接受但加上idx_user_rank联合索引后单用户查询稳定在 5 毫秒以内。4.3 用 JDBC 批量回写 MySQLfrom pyspark.sql.functions import row_number, col from pyspark.sql.window import Window results_with_rank results.withColumn( rank_no, row_number().over( Window.partitionBy(user_id).orderBy(col(score).desc()) ) ).filter(col(rank_no) 10) results_with_rank.repartition(4).write.jdbc( urljdbc:mysql://192.168.10.11:3306/restaurant?rewriteBatchedStatementstrue, tablerecommendations, modeoverwrite, properties{ user: spark_user, password: ******, driver: com.mysql.cj.jdbc.Driver, batchsize: 1000 } )这里有两个工程细节。第一repartition(4)限定了写库的并发连接数如果不做这一步Spark 默认的 200 个分区会同时打开 200 条 MySQL 连接很容易打爆数据库的 max_connections。第二rewriteBatchedStatementstrue让驱动把多条 INSERT 合并成一条多值语句提交写入吞吐量通常能提升 5 到 8 倍。modeoverwrite对 JDBC 的语义是 DROP TABLE 后重建重算过程中如果任务失败结果表会处于缺失状态。更稳的写法是先把结果写入recommendations_tmp成功后执行ALTER TABLE recommendations RENAME TO recommendations_bak和ALTER TABLE recommendations_tmp RENAME TO recommendations再做一次旧表清理。这样线上查询永远能读到完整数据。4.4 新订单如何尽早生效离线重算的频率通常是一天一次新下的订单最迟第二天早上生效。如果业务要求小时级更新可以在每小时用增量订单更新评分表并对部分用户做增量预测但 ALS 不支持流式增量训练只能把最近一天的新增行为拼进完整数据集重训。多数餐饮推荐场景日粒度已经足够因为用户一天打开菜单的次数有限推荐列表稳定性反而比新鲜度更重要。5. 推荐效果验证与调优的四个落地技巧这一章是真正拉开项目分数的地方也是业务上线前必须想清楚的部分。5.1 隐式反馈场景不要只看 RMSE显式评分可以用 RMSE 衡量预测误差但隐式反馈场景用户没有真实打分prediction与rating之间做 RMSE 没有语义价值。常见做法是把推荐当二分类评估真实购买过的(user_id, dish_id)对是正样本随机采样用户没买过的菜品作为负样本合并后对每对做模型打分计算 AUC。from pyspark.ml.evaluation import BinaryClassificationEvaluator # 正负样本合并后得到评估集 eval_df含 label 和 prediction 两列 evaluator BinaryClassificationEvaluator( labelCollabel, rawPredictionColprediction) auc evaluator.evaluate(eval_df)AUC 之外Top-N 命中率更贴近业务对测试集里每个用户统计他测试期真实购买的菜品有多少出现在模型给的 Top-10 列表里除以真实购买数得到 Recall10。这个指标可以直接写到项目答辩的 PPT 里。5.2 冷启动兜底coldStartStrategydrop会让新用户和新菜品拿不到推荐。常见兜底策略是新用户没进训练集时返回当日热销榜 Top-10新菜品没有隐因子时用同 category_id 下平均分最高的菜品占位。热销榜的 SQL 可以在order_items上按菜品聚合最近 7 天销量得出不需要 Spark 参与。5.3 数据倾斜和冷门菜完全无推荐餐饮数据里 20% 的热门菜往往占据 80% 的订单隐式 ALS 的置信度权重会让热门菜在推荐列表里反复出现挤压长尾。处理方式是控制单个菜品的订单量上限在构造评分时对buy_cnt做波峰截断比如超过 50 单的按 50 计算。另一个可尝试的方向是降低regParam让冷门菜能学到更贴近真实偏好的小数值隐因子。5.4 离线评估前检查这四件事训练数据时间窗口与测试集严格隔离implicitPrefs与评分构造方式匹配写库前确认rank_no已按 score 降序重排模型版本号写入recommendations.model_version方便回溯。这四个检查点全部通过推荐系统的离线链路才算真正闭合。本文还有配套的精品资源点击获取
返回列表