资讯动态

基于Python+Spark的汽车推荐系统:ALS协同过滤与调优

发布时间:2026/9/12 22:00:56 来源:尧图企业网站定制
简介基于PythonSpark的汽车推荐系统毕业设计资料包面向计算机相关专业在校学生与开发者适用于毕业设计、课程设计及项目初期立项参考内容围绕推荐系统构建与大数据处理流程展开。压缩包共29个文件涵盖Python爬虫脚本、Scala/Java源码、Markdown设计文档、多张界面与架构截图以及项目授权码总计7.99MB目录结构清晰便于按模块研读。目前已有127人学习下载项目经导师指导认可答辩评审分达95分代码均测试运行成功。除完整可运行的推荐系统源码外资料还包含汽车数据采集、Spark算子示例如ReduceByKeySort、Redis工具类JedisUtil及大屏可视化设计截图可帮助读者快速复现项目并理解数据清洗、特征计算与结果展示的完整链路。对需要提交高质量毕业设计或学习Spark推荐系统实践的读者而言这是一份高性价比的参考。1. 基于 PythonSpark 的汽车推荐系统先想清楚要解决什么汽车推荐系统这类大数据毕业设计难点不在算法本身而在“数据的真实感”和“计算过程的可视化”。同样的 ALS 模型用 Pandas 在本地跑 5 万条数据和用 PySpark 在集群上跑 200 万条数据时间和资源的说法完全不同。做这个题需要先明确要解决的关键问题面对低频、稀疏且高价值的汽车消费行为如何构建一张可训练的评分表如何在分布式环境下训练并调优协同过滤模型以及如何把结果落成可展示的推荐榜单。适合作这个题目的读者是已经懂一点 Python但还没有完整跑通 Spark MLlib 流程的大数据专业学生以及找工作前想拿一个完整推荐项目梳理算法和集群经验的初级工程师。2. PythonSpark 推荐引擎的选型原理为什么用协同过滤而不是深度学习汽车推荐的首选算法是协同过滤而不是深度学习模型。协同过滤需要的输入只有用户-汽车-评分三元组汽车信息表做冷启动兜底就够深度学习网络需要大规模连续行为序列和多模态特征毕设的数据量支撑不起来。Spark 集群的作用则是把几十万条行为数据和矩阵分解过程分到多个 executor 上并行算让“用了大数据框架”这件事在架构图和运行日志里都可以体现。2.1 汽车推荐场景的数据形态与选题边界在一开始把输入数据收敛成三张表。数据表关键字段在推荐中的用途用户行为表user_id, car_id, action, event_time, channel转换成评分是协同过滤的标签汽车信息表car_id, brand, price_level, fuel_type, body_type, seats构造内容特征解决冷启动用户信息表user_id, age, city, budget_level, family_size用户侧分析不对 ALS 输入做协同过滤建模时只需要前两张表。汽车信息表不在训练中使用而是在推荐结果生成后做业务过滤用户预算 15 万以内就不推荐 30 万以上的车型用户明确只看燃油车就把纯电车从榜单中去掉。数据集合中行为表的核心是评分汽车表的核心是白名单规则两者职责要分开别把品牌均价塞进评分里。数据量边界要在第一天就定下来。我一般用模拟行为生成器造 30 万到 50 万条记录覆盖 2 万用户、2000 辆车每个人行为 10 到 30 条。太多会拖慢调参太少又体现不出 Spark 的价值。真实数据集如果只有几千条建议适当补充隐式反馈否则 ALS 的收敛曲线会很难看。2.2 Spark 在汽车推荐系统里承担的角色从单机 DataFrame 到集群第二个选型问题是Spark 在这里到底承担什么如果只是read_csv后转 pandas 调 sklearn那答辩时被问到“集群跑在哪”会很难受。建议从第一步就保持 DataFrame 的分布式生命周期训练前不要collect()。from pyspark.sql import SparkSession spark ( SparkSession.builder .appName(CarRecommenderALS) .master(yarn) # 本地调试可临时改成 local[*] .config(spark.executor.memory, 2g) .config(spark.executor.cores, 2) .config(spark.sql.shuffle.partitions, 200) .getOrCreate() ) ratings_df ( spark.read.option(header, True) .csv(hdfs:///data/ratings.csv) ) ratings_df.printSchema()这里值得在论文里展开的是spark.sql.shuffle.partitions。它决定了 join、groupBy、窗口函数生成的 shuffle 分区数。推荐系统数据量不大200 个分区已经够用如果 executor 内存只有 2g分区数过大会导致每个分区调度成本大于计算成本。master(yarn)表示提交到 Hadoop YARN 集群先在本地启动一个local[*]session 把流程跑通再切 yarn排错成本最低。集群搭建时最少一主两从配置留在论文的实验环境小节。能说出这段“本地到集群”的切换过程比只贴 AUC 更有说服力。2.3 稀疏评分与冷启动ALS 的算法假设在哪里失效ALS 把用户和汽车映射成两组隐因子用点积拟合评分。这个模型成立的前提是每个用户和每辆车都参与了足够多的评分。汽车消费恰恰相反一个用户一年可能只在平台留下 5 条浏览记录很多冷门车型只有个位数打分ALS 对它们的预测会偏向全局均值。冷启动的解法不在模型里在数据里。第一把浏览、收藏、询价、下单映射成不同权重让数据密度增加第二对热门车和冷门车做分层抽样避免训练完全被头部车主导。下面代码把行为事件映射为评分from pyspark.sql import functions as F expr case action when detail then 1.0 when collect then 2.0 when inquiry then 3.0 when order then 5.0 else 0.5 end ratings_raw ( spark.read.option(header, True) .csv(hdfs:///data/user_behavior.csv) .withColumn(rating, F.expr(expr)) ) ratings_df ( ratings_raw .groupBy(user_id, car_id) .agg( F.max(rating).alias(rating), F.max(event_time).alias(event_time) ) )这里用max而不是sum是因为汽车属于重决策商品同一个用户连续查看同一辆车不代表意向翻倍只代表他还在考虑。用max保留最高行为级别能让标签更贴近真实购买漏斗。后面的 ALS 训练时implicitPrefs仍然设置成 False因为评分已经是 1 到 5 的浮点不是 0/1 隐式信号。提示模拟数据生成时行为时间要符合业务周期比如晚上和周末浏览多、工作日上午询价多否则按时间切分后测试集分布会和训练集明显不一致。3. 汽车推荐系统的数据准备用 PySpark 构建用户行为评分集上一章已经得到user_id, car_id, rating, event_time但直接拿去训练大概率会得到一个“看起来不错、实际不可解释”的结果。原因在于没有检查评分分布也没有设置合理的训练测试切分。这一章的处理过程也是答辩时“数据工程能力”的主要展示面。3.1 清洗评分数据剔除空值、超范围评分和超长行为窗口ALS 对输入数据里的异常值非常敏感。推荐系统多用离线评分跑出来的 RMSE 差 0.1 很可能就来自某条订单行为被错误记成了 10 分。先做一层基础清洗from pyspark.sql import Window from pyspark.sql import functions as F valid_rating_df ( ratings_df .filter(F.col(rating).isNotNull()) .filter(F.col(rating).between(1.0, 5.0)) .withColumn(rn, F.row_number().over( Window.partitionBy(user_id).orderBy(F.col(event_time).desc()) )) .filter(F.col(rn) 30) .drop(rn) .groupBy(user_id, car_id) .agg( F.max(rating).alias(rating), F.max(event_time).alias(event_time) ) )row_number按用户分区、按事件时间倒序排序只保留每个人最近的 30 条行为。这个窗口长度可以根据平均行为条数调整如果用户平均行为数是 15窗口设为 30 绰绰有余如果超过 50说明模拟数据里混进了异常用户。清洗后输出 DataFrame 只有四列这四列是冷启动和后续规则重排之前的唯一训练入口。3.2 按时间窗口切分训练集与验证集避免时间穿越训练集和测试集的切分方式直接影响评估指标的可信度。不要直接对整个 DataFrame 执行randomSplit([0.8, 0.2])那样同一个用户对同一辆车的记录会同时落在训练集和验证集中ALS 等于偷看了部分答案。更合理的做法是对每个用户按事件时间划分前 80% 做训练、后 20% 做验证。w Window.partitionBy(user_id).orderBy(F.col(event_time).asc()) split_df ( valid_rating_df .withColumn(rn, F.row_number().over(w)) .withColumn(total, F.count(*).over(Window.partitionBy(user_id))) ) train_df split_df.filter(F.col(rn) F.col(total) * 0.8).drop(rn, total) test_df split_df.filter(F.col(rn) F.col(total) * 0.8).drop(rn, total)这种切分有几个细节。一是event_time必须真实存在不能用随机数代替二是每个用户的rn会重新从 1 开始所以测试集里的交互对在训练集里从未出现过但用户 ID 本身在训练集里存在ALS 可以正常学习用户因子三是如果用户只有 1 到 3 条行为后 20% 可能是空集。这类用户可以直接过滤掉否则测试集里会出现大量只有训练没有测试的用户train_df train_df.join( test_df.select(user_id).distinct().withColumnRenamed(user_id, tu), train_df.user_id F.col(tu), inner ).drop(tu)这段 join 的作用是只保留那些“同时在训练集和测试集都有数据”的用户保证评估时每个用户都有至少一条待预测记录。我通常还会把行为总数低于 5 的用户在切分前直接过滤掉避免边界情况过度影响指标。3.3 检查评分分布确认 Spark 的 shuffle 是否失衡训练前用 SQL 看一眼数据总量和稀疏度。推荐系统调参大部分时间花在检查分布上而不在跑模型。SELECT COUNT(*) / COUNT(DISTINCT user_id) AS avg_actions_per_user, COUNT(DISTINCT car_id) AS car_count, COUNT(DISTINCT user_id) AS user_count FROM train_df;如果avg_actions_per_user低于 5说明行为窗口太短或者模拟数据太稀疏如果car_count低于 200说明汽车 SKU 太少模型容易把所有用户都推到同一个头部车型上。下面是一组可以作为参考的合理范围指标合理范围说明用户数1 万 - 10 万低于 1 万体现不出 Spark 优势汽车数500 - 5000汽车 SKU 比电商少一个量级人均行为数6 - 30小于 5 时需要补充隐式行为总评分条数10 万以上支撑 rank 20 以上的矩阵分解数据检查结果要单独存一张截图放进论文的数据分析章节这也是“详细文档”里最有含金量的一页。4. Spark ALS 模型训练与调参汽车推荐系统的核心参数表数据准备好之后模型本身并不复杂。Spark MLlib 的 ALS 在pyspark.ml.recommendation包里三列 DataFrame 直接丢进去就能训练。真正有信息量的是参数选择和评估口径。4.1 用 ALS 训练和预测显式评分与冷启动策略from pyspark.ml.recommendation import ALS from pyspark.ml.evaluation import RegressionEvaluator als ALS( userColuser_id, itemColcar_id, ratingColrating, rank20, maxIter15, regParam0.1, implicitPrefsFalse, coldStartStrategydrop, seed42 ) model als.fit(train_df) predictions model.transform(test_df) evaluator RegressionEvaluator( metricNamermse, labelColrating, predictionColprediction ) rmse evaluator.evaluate(predictions) print(fRMSE {rmse:.4f})coldStartStrategydrop的作用是丢弃那些在训练集中从未出现过的汽车或用户对应的预测值否则 transform 结果里会出现大量 nullRMSE 直接报错。显式评分场景下implicitPrefsFalse因为评分已经是 1 到 5 的浮点值如果数据源是点击 0/1 或者收藏 0/1才需要改成 True 并配合 alpha 参数。4.2 三个必调参数rank、maxIter 和 regParam毕业设计答辩时最容易被追问的就是“这些参数为什么这么设”。ALS 的参数不复杂但每个都有明确含义。参数尝试范围过大/过小的影响rank10 / 20 / 50过大会过拟合过小欠拟合maxIter10 / 15 / 20过大会让训练时间线性增加后期 loss 基本不再下降regParam0.01 / 0.1 / 0.5过小测试集 RMSE 高过大推荐结果趋同rank是隐因子个数。汽车数据集中汽车数只有几千rank 50 已经可以覆盖主要车型差异。regParam是正则化强度数据稀疏时建议从 0.1 开始调因为每个用户只有十几条行为模型很容易把训练集背下来。我自己的经验是先用固定maxIter15跑一次看 RMSE 和 loss再决定 rank 方向不要一上来就网格搜索否则每个参数组合都要跑完整轮训练。手动网格搜索的代码可以这样写from pyspark.ml.tuning import ParamGridBuilder, TrainValidationSplit param_grid ( ParamGridBuilder() .addGrid(als.rank, [10, 20, 50]) .addGrid(als.regParam, [0.01, 0.1, 0.5]) .build() ) tvs TrainValidationSplit( estimatorals, estimatorParamMapsparam_grid, evaluatorRegressionEvaluator( metricNamermse, labelColrating, predictionColprediction ), trainRatio0.8, parallelism2 ) tvs_model tvs.fit(train_df) best_model tvs_model.bestModel print(best_model.rank, best_model.parent.getRegParam())这里用TrainValidationSplit而不是CrossValidator是因为 ALS 训练本身较慢折数过多在集群上会成倍增加 task 数。毕设场景用一次切分足够把网格控制在 9 个组合以内。parallelism2表示同时运行两个参数组合可以根据 executor 数量适当调大。4.3 评估不能只看 RMSE还要看召回率和推荐覆盖率RMSE 衡量的是预测分和真实分的误差但推荐系统最终给用户看的是排序后的 Top-N。一个模型 RMSE 低有可能只是把大众车型的分数预测得准冷门车型永远排不上去。所以还要算两个业务指标Hit Rate 和覆盖率。recommendations best_model.recommendForAllUsers(10) rec_df ( recommendations .select(user_id, F.explode(recommendations).alias(rec)) .select(user_id, rec.car_id, rec.rating) ) hit_df rec_df.join(test_df, [user_id, car_id], inner) hit_users hit_df.select(user_id).distinct().count() hit_rate hit_users / test_df.select(user_id).distinct().count()hit_rate的含义是有多少比例的用户其推荐列表里至少命中了他后来真实交互的汽车。这个值比 RMSE 更贴近业务。覆盖率计算则直接看推荐列表覆盖了多少比例的汽车 SKUrec_car_cnt rec_df.select(car_id).distinct().count() total_car_cnt train_df.select(car_id).distinct().count() coverage rec_car_cnt / total_car_cnt覆盖率小于 30% 时说明推荐结果长期集中在少数头部车协同过滤没有发挥作用。此时应该降低 rank或者给冷门车型评分加一个小权重扰动。5. 让推荐结果可演示把 Spark 输出做成汽车大屏和答辩素材这一章把结果展示出来。最后一步不只是df.show()而是要把模型输出转换成可以直接上可视化大屏的 JSON。Spark 在这里的任务已经结束剩下的是把维度聚合的结果导出到前端。# 写出推荐明细供后续查询 rec_df.write.mode(overwrite).parquet(output/recommendations.parquet) # 品牌维度聚合 brand_df ( rec_df.join(car_info_df, car_id, left) .groupBy(brand) .agg(F.countDistinct(user_id).alias(rec_user_cnt)) .orderBy(F.col(rec_user_cnt).desc()) ) # 价格区间维度聚合 price_df ( rec_df.join(car_info_df, car_id, left) .groupBy(price_level) .agg(F.count(car_id).alias(rec_cnt)) .orderBy(price_level) ) data_for_frontend { brand_data: brand_df.toPandas().to_dict(orientrecords), price_data: price_df.toPandas().to_dict(orientrecords), }注意toPandas()只能在聚合结果已经很小的时候用。品牌数量只有几十个价格区间不到十个导出到本地方便 Flask 返回如果后端需要支撑在线请求应该把brand_df写回 MySQL 或者 Redis而不是每次查询都启动 Spark。线上演示时推荐列表还要叠加业务规则用户偏好燃油车时过滤掉纯电车型预算 15 万以内时过滤价格超过 20 万的车辆。这一步用 Spark SQL 写一个 WHERE 条件就够了放在rec_df之后执行逻辑清晰且好截图。答辩前记录三个证据Spark UI 里 ALS stage 的执行时间截图、RMSE 与覆盖率的变化曲线、以及 Top-N 推荐列表的对比样例。把explain()得到的训练日志存成文本放到论文实验部分比单独放一个混淆矩阵有力得多。本文还有配套的精品资源点击获取

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

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

免费获取报价