基于Python+Spark的汽车推荐系统:ALS协同过滤与调优
2026/9/12 22:00:51 网站建设 项目流程

简介:基于Python+Spark的汽车推荐系统毕业设计资料包,面向计算机相关专业在校学生与开发者,适用于毕业设计、课程设计及项目初期立项参考,内容围绕推荐系统构建与大数据处理流程展开。压缩包共29个文件,涵盖Python爬虫脚本、Scala/Java源码、Markdown设计文档、多张界面与架构截图以及项目授权码,总计7.99MB,目录结构清晰便于按模块研读。目前已有127人学习下载,项目经导师指导认可,答辩评审分达95分,代码均测试运行成功。除完整可运行的推荐系统源码外,资料还包含汽车数据采集、Spark算子示例(如ReduceByKeySort)、Redis工具类JedisUtil及大屏可视化设计截图,可帮助读者快速复现项目并理解数据清洗、特征计算与结果展示的完整链路。对需要提交高质量毕业设计或学习Spark推荐系统实践的读者而言,这是一份高性价比的参考。

1. 基于 Python+Spark 的汽车推荐系统,先想清楚要解决什么

汽车推荐系统这类大数据毕业设计,难点不在算法本身,而在“数据的真实感”和“计算过程的可视化”。同样的 ALS 模型,用 Pandas 在本地跑 5 万条数据和用 PySpark 在集群上跑 200 万条数据,时间和资源的说法完全不同。做这个题,需要先明确要解决的关键问题:面对低频、稀疏且高价值的汽车消费行为,如何构建一张可训练的评分表,如何在分布式环境下训练并调优协同过滤模型,以及如何把结果落成可展示的推荐榜单。适合作这个题目的读者,是已经懂一点 Python,但还没有完整跑通 Spark MLlib 流程的大数据专业学生,以及找工作前想拿一个完整推荐项目梳理算法和集群经验的初级工程师。

2. Python+Spark 推荐引擎的选型原理:为什么用协同过滤而不是深度学习

汽车推荐的首选算法是协同过滤,而不是深度学习模型。协同过滤需要的输入只有用户-汽车-评分三元组,汽车信息表做冷启动兜底就够;深度学习网络需要大规模连续行为序列和多模态特征,毕设的数据量支撑不起来。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( userCol="user_id", itemCol="car_id", ratingCol="rating", rank=20, maxIter=15, regParam=0.1, implicitPrefs=False, coldStartStrategy="drop", seed=42 ) model = als.fit(train_df) predictions = model.transform(test_df) evaluator = RegressionEvaluator( metricName="rmse", labelCol="rating", predictionCol="prediction" ) rmse = evaluator.evaluate(predictions) print(f"RMSE = {rmse:.4f}")

coldStartStrategy="drop"的作用是丢弃那些在训练集中从未出现过的汽车或用户对应的预测值,否则 transform 结果里会出现大量 null,RMSE 直接报错。显式评分场景下implicitPrefs=False,因为评分已经是 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 开始调,因为每个用户只有十几条行为,模型很容易把训练集背下来。我自己的经验是:先用固定maxIter=15跑一次,看 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( estimator=als, estimatorParamMaps=param_grid, evaluator=RegressionEvaluator( metricName="rmse", labelCol="rating", predictionCol="prediction" ), trainRatio=0.8, parallelism=2 ) tvs_model = tvs.fit(train_df) best_model = tvs_model.bestModel print(best_model.rank, best_model.parent.getRegParam())

这里用TrainValidationSplit而不是CrossValidator,是因为 ALS 训练本身较慢,折数过多在集群上会成倍增加 task 数。毕设场景用一次切分足够,把网格控制在 9 个组合以内。parallelism=2表示同时运行两个参数组合,可以根据 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 更贴近业务。覆盖率计算则直接看推荐列表覆盖了多少比例的汽车 SKU:

rec_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(orient="records"), "price_data": price_df.toPandas().to_dict(orient="records"), }

注意toPandas()只能在聚合结果已经很小的时候用。品牌数量只有几十个,价格区间不到十个,导出到本地方便 Flask 返回;如果后端需要支撑在线请求,应该把brand_df写回 MySQL 或者 Redis,而不是每次查询都启动 Spark。

线上演示时,推荐列表还要叠加业务规则:用户偏好燃油车时,过滤掉纯电车型;预算 15 万以内时,过滤价格超过 20 万的车辆。这一步用 Spark SQL 写一个 WHERE 条件就够了,放在rec_df之后执行,逻辑清晰且好截图。

答辩前记录三个证据:Spark UI 里 ALS stage 的执行时间截图、RMSE 与覆盖率的变化曲线、以及 Top-N 推荐列表的对比样例。把explain()得到的训练日志存成文本,放到论文实验部分,比单独放一个混淆矩阵有力得多。

本文还有配套的精品资源,点击获取

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询