从去年开始我一直在做汽车销售方向的数字化项目,当时接到的需求很直接:公司旗下几十家4S门店的客户看车数据全在系统里,但是销售顾问给客户打电话推荐车型基本靠感觉,管理层看经营情况也只知道总数,不知道品牌结构、价格带趋势、客户流失风险。于是我们落地了一套以SpringBoot为应用层、Spark为计算引擎的汽车销售推荐与大数据分析系统,这套方案一直跑到现在,期间踩了不少坑,也提炼出了不少可以复用的经验。如果你是Java出身、想往大数据方向靠,或者正在做汽车、房产这类低频高价值行业的推荐系统,这篇文章应该能帮你少走很多弯路。
1. 项目定位与技术选型:为什么是SpringBoot加Spark这套组合
1.1 汽车销售场景下的推荐需求,和电商推荐根本不是一回事
做推荐系统前,我们先把业务场景捋了一遍。汽车消费和买衣服、点外卖差别很大,它是典型的长决策链、低频高价值消费,一个用户从看车到成交可能经历几周到几个月,行为路径往往是:线上浏览参数、收藏备选车型、到店询价、试驾、最终下单。这意味着我们不能简单照搬电商那套“点击了A就推荐B”的逻辑,而要看用户当前处于决策漏斗的哪一层,推送对应阶段的内容。
当时系统里已经积累了大约50万注册用户、2000多个在售车型SKU,日新增行为日志在几十万条量级。这个量级说实话不算特别大,但问题是行为数据极其稀疏:大多数用户只产生几次浏览行为就离开了,真正走到询价、试驾的比例非常低。如果直接在SpringBoot应用里用内存计算做推荐,数据量上来后光是JOIN和排序就会把服务拖垮,更别提还要跑模型训练。
所以我们的定位很清晰:SpringBoot负责对外服务、业务编排、接口输出,Spark负责离线批量计算、特征加工和模型训练,两边通过数据存储层衔接。这样应用层保持轻量,计算层可以独立扩容。
1.2 系统架构分层和数据流向
整体架构上我们分成了四层,每一层的职责都比较清晰:
- 应用层:SpringBoot提供REST接口,包括推荐接口、经营分析报表接口、后台管理接口,同时管理用户、车型、订单等业务数据。
- 存储层:MySQL存业务主数据,Redis做缓存,HDFS存行为日志原始文件,分析结果会写回MySQL或Redis供应用层读取。
- 计算层:Spark负责每天凌晨跑批,完成数据清洗、特征工程、ALS模型训练、销售指标聚合计算。
- 展示层:前端用Vue和ECharts展示看板,App端和顾问端通过接口拉取推荐结果。
数据流大概是:埋点日志进Kafka,定时落HDFS,同时业务库的订单、车型数据通过DataX同步到数仓。每天凌晨Spark读HDFS和业务库快照,产出两类结果:一类是每个用户的车型推荐列表,写入Redis给接口用;另一类是销售KPI聚合结果、漏斗转化指标,写入MySQL给报表系统用。
1.3 为什么不用纯Java内存计算,也不用Python加PySpark
项目立项时团队内部有过方案分歧。一部分同事觉得数据量才百万级,直接用SpringBoot内存算也行;另一部分想上Python和PySpark,说生态好写起来快。最后我们两个都没选,理由如下:
纯Java内存计算的问题在于,刚开始数据量小确实能跑,但一旦要加特征维度、做交叉统计、跑协同过滤矩阵分解,单机内存和CPU都扛不住,代码里全是并发和线程池的坑,后期维护成本极高。而Python加PySpark的问题在于团队背景不一样——我们主力是Java工程师,Python只是辅助写脚本,如果核心计算链路用PySpark,出了线上问题没人敢接。
最终选型是Java写Spark作业,用Spark的Java API完成训练和计算,再由SpringBoot直接调用计算结果。这套方案的好处是整个代码库统一在Java技术栈,运维、排障、人员招聘都方便,性能上也没损失多少。
2. 数据建模与特征工程:推荐系统能不能出效果,七成看这里
2.1 数据源和埋点设计
推荐系统最怕的就是“数据都还没采好,就急着上模型”。我们第一版踩过这个坑,当时行为日志只记录了“用户点击了哪款车”,没有上下文信息(页面位置、停留时长、是否来自推荐位),导致后面做特征工程时很多指标算不出来。后来重新梳理了埋点,核心数据源分三块:
- 用户基础表:用户id、性别、年龄段、所在城市、购车预算区间、增换购状态。
- 车型表:车型id、品牌、指导价、车辆级别(紧凑型/中型/中大型/SUV/MPV)、能源类型(燃油/纯电/混动)、上市时间、标签(家用、运动、商务等)。
- 行为表:行为类型(浏览、收藏、询价、试驾、下单、成交)、行为时间、行为来源(自然流量/推荐位/搜索)、停留时长、是否异常点击。
埋点字段虽然不多,但足够支撑后续的行为评分和漏斗分析。这里有个经验:宁可多埋几个字段,也不要等模型需要了再补采集,补数据是最痛苦的事。
2.2 行为评分与隐式反馈处理
汽车销售场景下几乎没有用户主动打分“我喜欢这款车”,所以我们面对的是典型的隐式反馈数据。隐式反馈的特点是用户没做某事不代表不喜欢,行为频次低也不代表意向弱,必须设计一套合理的评分映射。
我们当时的评分规则是:浏览一次1分,收藏3分,询价5分,试驾8分,成交10分。然后按用户和车型聚合,一个用户对一辆车的总行为分就是训练用的rating。另外我们还对行为时间做了衰减:30天内的行为权重为1,30到90天权重衰减到0.6,90天以上衰减到0.3。原因很简单,用户三个月前看过的车,现在可能已经提车或者换目标了,不能和新行为同等对待。
还要注意数据清洗。当时日志里有不少爬虫和无效点击,我们通过两个规则过滤:同一用户对同一车型单日点击超过20次直接判定异常,只保留最高的3次;来源标记为“内部测试”的设备id全部剔除。不过滤这些脏数据,模型训练出来会有明显的偏移。
2.3 用Spark把原始数据加工成训练集
数据准备阶段的Spark作业,核心逻辑就是从HDFS读取埋点日志,join业务库的车型表和用户表,按评分规则生成“用户—车型—评分”三元组。这一步看起来简单,实际写的时候要特别注意数据倾斜。
我们当时的处理方式是读完原始日志后先做一次按行为日期的分区裁剪,再按用户id做桶聚合。因为有些热门车型的行为量是冷门车型的几十倍,如果不加预聚合,后面按车型维度做统计时容易把某个executor打爆。核心代码大概是这样的:
SparkSession spark = SparkSession.builder() .appName("CarBehaviorETL") .config("spark.sql.shuffle.partitions", "200") .enableHiveSupport() .getOrCreate(); Dataset<Row> logs = spark.read().parquet("hdfs://nameservice1/data/behavior_logs"); Dataset<Row> cars = spark.read().jdbc(jdbcUrl, "car_model", props); Dataset<Row> users = spark.read().jdbc(jdbcUrl, "user_info", props); Dataset<Row> joined = logs.join(cars, logs.col("car_id").equalTo(cars.col("id"))) .join(users, logs.col("user_id").equalTo(users.col("id"))) .filter("is_test = 0"); Dataset<Row> training = joined .select( logs.col("user_id").cast("long"), logs.col("car_id").cast("long"), scoreExpr().as("rating") ) .groupBy("user_id", "car_id") .agg(functions.sum("rating").as("rating")) .filter("rating > 0"); training.write().mode("overwrite").parquet("hdfs://nameservice1/data/training_set");生成好的三元组数据就是后面ALS训练的输入,我们一般会保留最近三个月的有效行为,太旧的数据对当前推荐没有帮助,反而引入噪声。
3. 推荐引擎实现:ALS协同过滤加冷启动兜底
3.1 为什么选ALS而不是UserCF或ItemCF
主流的协同过滤思路有三种:基于用户的UserCF、基于物品的ItemCF、基于矩阵分解的ALS。在汽车这种稀疏、低频的场景里,我们最终选了ALS。
UserCF的思路是找到和你行为相似的用户,推荐他们看过的车。问题在于汽车用户行为太稀,能算出的相似用户非常有限,推荐结果容易坍缩成热门车。ItemCF的思路是推荐和你浏览过的车相似的车,效果稍好,但同样受限于“这个用户的历史行为实在太少”。而且这两种算法在Spark里实现起来虽然不复杂,但用户维度扩展后计算量增长很快。
ALS矩阵分解则把用户和车型映射到同一个低维隐向量空间:用户向量代表用户的偏好画像,车型向量代表车型的属性特征,两者点积就是对用户-车型匹配度的预测。Spark MLlib里直接封装了ALS算法,分布式训练稳定性很高,对稀疏矩阵的处理也比较友好。
有人可能会问,汽车推荐要不要上深度学习模型、上知识图谱?我们的判断是没必要。在数据量只有百万级、行为极度稀疏的情况下,复杂模型很容易过拟合,而且解释性差,不好向业务方交代。ALS这个级别的模型已经能覆盖大多数场景,先把基础命中率做上来,再谈花活。
3.2 Spark MLlib ALS训练的关键参数与实现
我们用的是Spark MLlib里的ALS实现,训练代码用Java写,关键代码如下:
import org.apache.spark.ml.recommendation.ALS; import org.apache.spark.ml.recommendation.ALSModel; import org.apache.spark.sql.Dataset; import org.apache.spark.sql.Row; import org.apache.spark.sql.SparkSession; SparkSession spark = SparkSession.builder() .appName("CarALSRec") .config("spark.serializer", "org.apache.spark.serializer.KryoSerializer") .getOrCreate(); Dataset<Row> training = spark.read().parquet("hdfs://nameservice1/data/training_set"); ALS als = new ALS() .setMaxIter(10) .setRank(20) .setRegParam(0.1) .setUserCol("user_id") .setItemCol("car_id") .setRatingCol("rating") .setColdStartStrategy("drop"); ALSModel model = als.fit(training); model.write().overwrite().save("hdfs://nameservice1/models/car_als_" + dateStr);参数选择方面,我们调过不少组合,最后稳定在rank=20、maxIter=10、regParam=0.1,这个组合在验证集上的RMSE最低。其中rank表示隐向量维度,太大会导致过拟合和内存压力,太小则表达不了用户和车型的复杂特征;regParam是正则化系数,用来防止过拟合,可以理解为对模型复杂度的惩罚,数据越稀疏,这个值通常要调大一些。
还有一点必须提:MLlib的ALS默认setColdStartStrategy是nan,意思是如果新用户或新车型没有训练数据,预测结果直接返回NaN,这在线上是会出事故的。我们显式设置成drop,并且要求训练和预测前过滤掉没有向量表示的用户。
3.3 冷启动问题的三个兜底策略
ALS训练出的模型只能覆盖“行为丰富”的用户,真正上线时你会发现大量新注册用户和只点过一次车的人根本没有预测结果。这时候如果直接把空列表返回给前端,用户体验会很差,业务方也会质疑推荐系统的价值。
我们当时设计了三个兜底策略,按优先级从高到低:
- 新用户兜底:直接返回全站车型综合热度榜TopN,热度分参考浏览、收藏、询价、成交的加权汇总,保证推荐结果不冷场。
- 低行为用户兜底:如果该用户历史行为少于3次,优先召回该用户所在城市销量靠前、与其浏览过的车型同级别或同能源类型的车。
- 有预测结果但置信度低:ALS推荐结果里评分相近的多个车型都保留,不做硬截断,方便用户在推荐池中切换。
这三个策略不是拍脑袋定的,核心逻辑是:在用户数据足够支撑个性化之前,先给一个“不犯错”的推荐;当用户行为慢慢积累后,再逐步加大ALS结果的权重。我们在Redis里存的推荐结果会带上来源标签,比如rec:user:123对应的list里每辆车后面标注ALS、HOT或RULE,这样出问题的时候可以快速定位是哪个渠道出了问题。
4. 数据分析模块:从销售看板到用户分群
4.1 销售KPI指标拆解与Spark SQL实现
光有推荐还不够,领导层更关心的是经营分析。我们做了一套销售分析看板,核心指标包括:总销量及环比同比、品牌销量份额、动力类型占比、价格带分布、区域销售热度、库存与成交价背离度。
这些指标如果用SpringBoot直接查MySQL算,光是一次全量聚合就能把业务库拖到慢查询告警。我们改成Spark SQL每晚跑批,先把订单流水、车型档案、门店信息同步到数仓,再用宽表方式做聚合,最后把结果写回MySQL报表表。聚合SQL大概是这种风格:
SELECT DATE_FORMAT(order_date, 'yyyy-MM') AS month, brand, COUNT(*) AS sales_cnt, SUM(deal_price) AS sales_amount, ROUND(AVG(deal_price), 2) AS avg_deal_price FROM dwd_car_order WHERE order_date >= '2024-01-01' GROUP BY DATE_FORMAT(order_date, 'yyyy-MM'), brand;环比和同比我们用窗口函数算,比在SpringBoot里逐个月查再比要高效得多:
SELECT month, brand, sales_cnt, LAG(sales_cnt, 1) OVER (PARTITION BY brand ORDER BY month) AS prev_month_sales, ROUND((sales_cnt - LAG(sales_cnt, 1) OVER (PARTITION BY brand ORDER BY month)) * 100.0 / LAG(sales_cnt, 1) OVER (PARTITION BY brand ORDER BY month), 2) AS mom_rate FROM brand_month_sales;这里还有个体会:窗口函数在Spark里的支持很完善,但要注意PARTITION BY的字段基数,如果维度值特别多,shuffle成本会很高。我们后来对品牌这类低基数字段做了广播变量优化,跑批时间从40分钟降到20分钟以内。
4.2 用户分群与RFM模型
销售数据的第二个重点是用户分群。管理层想知道的不是“用户总数有多少”,而是“哪些用户是即将流失的高价值客户”“哪些用户值得优先跟进”,所以我们基于RFM模型做了分层。
RFM是指最近一次消费时间(Recency)、消费频率(Frequency)、消费金额(Monetary)。我们用Spark SQL从订单表算出每个用户的三个指标,再分别按分位数打1到5分,综合得出用户价值层级:
SELECT user_id, ROUND(SUM(deal_price), 2) AS total_amount, COUNT(DISTINCT order_id) AS order_cnt, DATEDIFF(CURRENT_DATE(), MAX(order_date)) AS days_since_last_order FROM dwd_car_order GROUP BY user_id;分层后的用户群体分成了高价值忠诚用户、成长型用户、流失预警用户、沉默用户四类。流失预警用户会进入运营任务池,由销售顾问优先跟进;高价值忠诚用户则会被打上标签,后续做换购推荐时权重更高。这里有个冷知识:在做汽车换购推荐时,RFM分层比单纯的行为协同过滤更有效,因为换购决策更多依赖历史消费能力和品牌忠诚度,而不是最近看了什么车。
4.3 分析结果怎么给前端用
数据分析结果最终要落到界面上。我们的做法是:Spark跑批完成后,把聚合结果写入MySQL的报表专用表,同时清理C端报表接口的Redis缓存,保证前端ECharts拉到的数据是最新一批。技术上SpringBoot这边就是普通的接口查询,没啥特殊。
比较值得说的是指标口径的统一。最开始各个业务方对“销量”的定义都不一样:有人认为是订单支付成功算销量,有人认为是开票算销量,还有人认为交车才算。我们花了很大力气把核心指标口径固化在数仓层,所有报表统一从dwd_car_order表取数,彻底消掉了“同一个数,两套报表数值不一致”的扯皮问题。这个经验对任何带数据分析的系统都适用。
5. SpringBoot和Spark集成的工程细节
5.1 两种集成方式,别选错
集成Spark和SpringBoot是很多人在项目初期会纠结的问题。我们其实试过两种方式,各有适用场景:
第一种是在SpringBoot进程内启动SparkSession,用local模式跑计算。开发阶段我特别推荐这种方式,因为调试方便,断点能直接打进RDD算子。但生产环境千万不要这么做,原因有三个:一是Spark的executor内存管理和SpringBoot的JVM堆容易打架,二是Spark作业长任务会占用应用进程的线程资源,三是应用重启会导致正在跑的Spark作业直接丢失。
第二种是独立Spark作业jar包,通过spark-submit提交到集群,用调度框架定时触发。我们生产环境用的就是这种。SpringBoot侧需要触发跑批时,就用ProcessBuilder调用spark-submit命令,把参数传进去。示例:
spark-submit \ --class com.company.recommend.ALSRecommendJob \ --master yarn \ --deploy-mode cluster \ --executor-memory 4g \ --executor-cores 2 \ --num-executors 10 \ --conf spark.sql.shuffle.partitions=200 \ hdfs://nameservice1/jars/car-recommend-1.0.jar \ --date 2024-06-01用这种提交方式,Spark作业和应用服务彻底隔离,跑批失败也不影响在线接口,只需要在调度平台配置重试和告警。
5.2 推荐结果接口缓存与降级设计
推荐结果接口是面向C端的高频接口,不可能每次请求都去查MySQL或重新算一遍推荐,必须做缓存。我们是这么设计的:
推荐结果每天凌晨由Spark跑批生成,写入Redis,接口层先查Redis,命中直接返回;没命中再查MySQL里的“上一版推荐结果”表,保证服务不裸奔。Redis的key结构是rec:user:{userId},value用JSON存一个推荐列表,每项包含车型id、推荐分、来源标签。列表长度控制在20个,前端首屏只展示10个,超过部分作为“换一批”的数据源。
缓存设置TTL为48小时,也就是最多两天旧数据。这个时间窗口对低频汽车消费刚好,用户今天看到推荐的车,明天也还适用,没必要像电商那样按小时更新。
5.3 定时调度与增量更新机制
算法模型不能只跑一次,需要定期更新。我们的调度策略是:每天凌晨2点执行全量ETL和模型训练,先生成用户行为快照和车型特征快照,再训练ALS模型,最后生成推荐结果并写缓存。白天如果业务库发生大额订单或试驾行为,通过消息队列触发轻量级增量更新,只更新受影响用户和热门车型的推荐结果。
这里有个重要的取舍:全量模型每天训练一次,但模型文件并不代表实时性,因为用户行为一直在产生。为了让推荐结果看起来“不那么滞后”,我们在生成推荐列表时不是直接输出ALS排名前N,而是把当天新增行为加进去做二次重排:如果用户今天刚收藏了某款车,这款车的推荐排名会强制提到前3。这个规则成本极低,但业务反馈“推荐终于像懂了我在看什么”。
6. 实际踩坑与性能调优实录
6.1 ALS训练数据稀疏引起的模型坍缩
第一次把ALS模型跑上线,我们发现推荐的车型集中在少数几款热门车上,基本等于热门榜,业务方很不满意。排查后发现原因有两个:一是很多用户的训练数据太少,二三行为就进了训练集,模型学到的向量没有区分度;二是评分聚合后分布极不均衡,热门车型的评分总和远大于长尾车型。
解决办法:一是过滤掉行为数少于3条的用户,这个阈值不能定太高,否则有效用户数太少;二是对热门车型做降采样,限制单个车型最多贡献的行为样本比例;三是调大正则化参数regParam,从0.01调到0.1,抑制过拟合。调完之后虽然离线RMSE略有上升,但线上推荐列表的多样性明显改善,长尾车型的曝光量提升不少。
6.2 Executor OOM和Shuffle压力
跑批作业最常遇到的就是executor OOM。我们第一版spark-submit只给了executor-memory 2g,结果每天跑的ETL作业在shuffle阶段频繁报OOM,日志里全是Container killed by YARN for exceeding memory limits。
排查后确认问题出在行为日志和车型表JOIN时,车型表虽然不大,但默认要复制到每个executor参与shuffle,造成大量网络和内存开销。优化方案是:把车型表、用户基础表这类小维度表做成广播变量(spark.sql.autoBroadcastJoinThreshold调大到20MB),避免shuffle;同时把executor内存调到4g、增加executor数量到10个,并开启Kryo序列化减少内存占用。调整后同一份作业执行时间从50分钟降到20分钟,OOM基本消失。
6.3 数据倾斜的定位与处理
另一个高频问题就是数据倾斜,特别是在按品牌、按门店聚合的统计作业里。现象是某个executor跑得特别慢,整个Stage卡住,其他executor都在空转等它。我们用Spark UI看到某个task处理的数据量是其他task的数十倍,基本可以确认是倾斜。
处理方式分两种:如果是大表和小表JOIN导致的倾斜,优先用小表广播;如果聚合Key本身分布不均(比如某品牌销量占50%),就对Key加随机前缀打散到多个task再聚合,或者用repartition重新调整分区。我们当时在品牌维度聚合时用了加盐思路:先给品牌字段加随机后缀拆成多组,算完再合并。
6.4 推荐效果评估的线上指标
最后说说怎么评估推荐系统到底有没有用。离线我们看RMSE,线上我们真正关心的是推荐位的点击率、推荐位带来的询价转化率、以及推荐位贡献的成交占比。在我们的数据里,推荐位贡献的PV不到全站10%,但带来的询价线索占比稳定在25%上下,这才是推荐系统的核心价值——不在于流量分发,而在于把高意向的用户提前捞出并推给销售跟进。
我个人的体会是,做这种垂直领域推荐系统,不必一上来就追最新的大模型、图神经网络,先把ALS、规则兜底、指标看板这套基础盘跑通,业务价值已经很明显了。再有就是在项目里多留一手可解释性:每次推荐都要能说清楚“为什么推荐这辆车给这个用户”,这在给销售顾问做辅助时特别重要——顾问跟客户聊车的时候需要理由,光丢个推荐结果过去是没有说服力的。