ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

SpringBoot+Spark打造汽车销售推荐系统:从协同过滤到冷启动实践

SpringBoot+Spark打造汽车销售推荐系统:从协同过滤到冷启动实践 从去年开始我一直在做汽车销售方向的数字化项目当时接到的需求很直接公司旗下几十家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天内的行为权重为130到90天权重衰减到0.690天以上衰减到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(); DatasetRow logs spark.read().parquet(hdfs://nameservice1/data/behavior_logs); DatasetRow cars spark.read().jdbc(jdbcUrl, car_model, props); DatasetRow users spark.read().jdbc(jdbcUrl, user_info, props); DatasetRow 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); DatasetRow 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(); DatasetRow 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);参数选择方面我们调过不少组合最后稳定在rank20、maxIter10、regParam0.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.partitions200 \ 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、规则兜底、指标看板这套基础盘跑通业务价值已经很明显了。再有就是在项目里多留一手可解释性每次推荐都要能说清楚“为什么推荐这辆车给这个用户”这在给销售顾问做辅助时特别重要——顾问跟客户聊车的时候需要理由光丢个推荐结果过去是没有说服力的。
返回列表