基于SpringBoot与Spark的汽车销售推荐系统与大数据分析实践
1. 项目背景与整体设计思路1.1 为什么要做这个系统先说一个比较现实的问题汽车4S店和销售平台手里的数据量其实非常庞大每台车的配置参数、用户浏览记录、试驾反馈、成交价格、售后记录这些数据每天都会产生上万条。但很多门店对这些数据的使用方式还停留在月底拉个Excel表看看销量的阶段根本谈不上数据分析更别说推荐系统。我去年接手了一个汽车销售平台的改造项目客户的核心诉求有两点一是想搞清楚哪些用户最可能买哪款车二是想盘点库存和销量之间的关联规律。传统做法是用业务数据库直接跑统计SQL但数据量一旦到了千万级别关联查询经常几十秒甚至几分钟才出结果体验非常差。后来我采用了SpringBoot Spark的组合把推荐和数据分析拆成两条链路才彻底解决了性能瓶颈。这篇文章我会把完整的方案讲清楚包括Spark集群怎么搭、推荐算法怎么实现、SpringBoot怎么和Spark配合、以及我实际踩过的坑。适合手里有汽车销售场景、想做推荐系统或者大数据分析的Java工程师参考也适合正在做毕业设计或课程项目的同学借鉴思路。1.2 技术选型背后的考量当初选型的时候我其实纠结过几个方案纯SpringBoot MySQL Redis简单但处理不了海量数据的复杂计算Flink实时流处理能力强但对团队的运维要求太高杀鸡用牛刀Spark SpringBoot既能做批处理分析又能通过MLlib实现推荐算法生态成熟Java和Scala都能写最终选了第四条路。Spark的最大优势是内存计算对汽车销售这种白天积累数据、晚上统一分析的场景非常契合。我只需要每天凌晨用Spark跑一次离线分析任务生成推荐结果和统计数据白天SpringBoot服务只要读取预计算结果即可QPS完全扛得住。这里有个关键认知推荐系统不一定要实时计算。对于汽车这种低频、高客单价的商品用户一天内不会反复浏览几百次离线计算定时更新的策略已经足够而且实现成本低了一个量级。1.3 系统架构一览整个系统的数据流是这样的业务端SpringBoot应用产生用户行为日志和订单数据写入MySQL同时通过AOP切面记录用户浏览、收藏、试驾等行为日志发送到KafkaSpark定时任务从Kafka拉取日志经过清洗后存入HDFSSpark MLlib跑协同过滤和商品相似度计算把结果写回MySQL和RedisSpringBoot通过REST API对外提供推荐接口和分析报表接口。这套架构的好处是查询和计算分离白天业务高峰期Spark不会抢数据库资源推荐结果预计算好接口响应时间稳定在10ms以内数据和逻辑解耦后续想换算法或者加数据源都不需要动主服务2. 数据层设计与预处理2.1 数据来源和埋点方案做推荐系统最怕的就是数据质量差。我见过太多项目连最基本的用户行为日志都没埋全最后算法跑出来的结果全是噪声。汽车销售场景下我建议至少采集以下几类数据用户基础信息性别、年龄段、所在城市、消费能力等级后台注册的信息浏览行为车牌型号、浏览时长、访问来源页、浏览深度互动行为收藏、对比、询价、预约试驾订单数据最终成交车型、成交价格、购买方式埋点方案我用的是SpringBoot拦截器加AOP的方式。用户每次点击车型详情、发起试驾预约、收藏车型时都会异步发送一条消息到Kafka。这里注意别用同步发送否则用户的点击请求会被消息队列拖慢我一开始犯过这个错误后面优化成异步发送后接口耗时从80ms降到了15ms。2.2 用户行为日志表设计清洗后的数据最终存储结构大概是这样的字段名类型说明userIdString用户IDcarModelIdString车型IDbehaviorTypeInt行为类型1浏览 2收藏 3对比 4询价 5试驾behaviorTimeTimestamp行为发生时间pageDurationInt停留时长秒channelString渠道来源设计时有两个注意点行为类型必须用数字枚举存字符串会浪费空间而且Spark处理时需要额外转换停留时长必须记录这是判断用户兴趣度的关键指标。如果用户只看了5秒就关掉这条浏览记录的权重应该很低如果看了一分钟以上说明用户确实感兴趣2.3 Spark数据清洗要点日志数据进HDFS之前我用Spark做了三步清洗第一步去重。因为Kafka可能重复投递消息清洗时要按userId carModelId behaviorTime behaviorType做去重组合键。去重我用的是dropDuplicates方法在Spark 2.4以上版本配合sort合并小文件效果很好。第二步过滤异常数据。比如停留时间超过2小时的记录正常用户不可能盯着一辆车看两小时大概率是挂了后台没关页面、凌晨刷量的爬虫记录等等。这一步不能省否则推荐结果会被异常数据带偏。第三步数据脱敏。用户手机号、身份证等敏感字段只保留部分位数虽然内网项目不一定需要但养成好习惯总没错。3. 推荐系统核心算法实现3.1 算法选型ALS协同过滤汽车销售推荐场景我最终选择了Spark MLlib中的ALS交替最小二乘法协同过滤算法。原因有三点ALS适合稀疏矩阵。汽车车型可能有两三百款用户最多浏览过其中十几种这个稀疏度在协同过滤里比较常见ALS的隐式反馈处理能力能很好适应。ALS训练速度快。在Spark集群上设置合理的迭代次数和正则化参数后跑30万用户对200款车型的评分矩阵训练时间可以控制在几分钟内。ALS有成熟的调参经验。不像深度学习模型需要反复摸索网络结构ALS的核心参数就三个rank、iterations、lambda理解成本低也好调整。3.2 ALS模型的数学原理ALS的核心思想是把用户-物品评分矩阵分解为两个低维矩阵的乘积。假设有m个用户n个物品评分矩阵R是一个m×n的矩阵。ALS把它分解为R ≈ U × V其中U是m×k的用户特征矩阵V是n×k的物品特征矩阵k就是rank参数表示潜在特征的个数。说白了就是给每个用户和每辆车都学一个k维的隐形向量这个向量捕捉了用户的偏好特征和车辆的属性特征。算法通过迭代求解来逼近真正的R。每次迭代时先固定物品矩阵V求解用户矩阵U的最优解再固定用户矩阵U求解物品矩阵V的最优解。交替进行直到收敛。这里我补充一个实操理解k值到底取多少合适。如果k太小模型拟合能力不足推荐结果不够准确k太大虽然训练集上的误差小但容易过拟合而且计算量成倍增加。我实测下来汽车销售场景k取值在30到50之间比较合适。3.3 Spark实现代码下面是核心的ALS训练代码我用的是Spark 2.4的MLlibimport org.apache.spark.ml.recommendation.ALS import org.apache.spark.sql.SparkSession // 创建SparkSession设置内存和并行度 val spark SparkSession.builder() .appName(CarSalesRecommendation) .master(yarn) .config(spark.sql.shuffle.partitions, 100) .getOrCreate() // 读取清洗后的行为数据 val ratingData spark.sql( SELECT userId, carModelId, behaviorScore as rating FROM cleaned_behavior_log WHERE behaviorDate 2024-01-01 ) // 划分训练集和测试集 val Array(trainingData, testData) ratingData.randomSplit(Array(0.8, 0.2)) // 配置ALS模型 val als new ALS() .setRank(40) // 潜在特征个数 .setMaxIter(15) // 最大迭代次数 .setRegParam(0.05) // 正则化参数 .setUserCol(userId) .setItemCol(carModelId) .setRatingCol(rating) .setColdStartStrategy(drop) // 冷启动时丢弃无法预测的评分 // 训练模型 val model als.fit(trainingData) // 预测测试集评分评估模型效果 val predictions model.transform(testData) import org.apache.spark.ml.evaluation.RegressionEvaluator val evaluator new RegressionEvaluator() .setMetricName(rmse) .setLabelCol(rating) .setPredictionCol(prediction) val rmse evaluator.evaluate(predictions) println(sRoot-mean-square error $rmse)模型训练完成后用model.recommendForAllUsers(10)可以给每个用户生成Top10推荐列表。输出结果直接写到MySQL的推荐结果表里。3.4 评分权重的设定用户没有给汽车打过评分所以我们需要把行为日志转化为评分。这个转化逻辑直接决定了模型的效果。我用的权重方案是浏览1分收藏3分对比5分询价8分试驾10分。然后按时间衰减处理最近一周的行为权重为1.0一周到一个月之间的权重为0.6一个月以上的权重为0.3。这样最近的行为对推荐的贡献更大。具体代码实现public Integer calculateScore(BehaviorLog log) { int baseScore; switch (log.getBehaviorType()) { case 1: baseScore 1; break; // 浏览 case 2: baseScore 3; break; // 收藏 case 3: baseScore 5; break; // 对比 case 4: baseScore 8; break; // 询价 case 5: baseScore 10; break; // 试驾 default: baseScore 0; } long daysDiff Duration.between(log.getBehaviorTime(), LocalDate.now()).toDays(); double timeDecay daysDiff 7 ? 1.0 : (daysDiff 30 ? 0.6 : 0.3); return (int) Math.round(baseScore * timeDecay); }提示权重值不是拍脑袋定的要结合业务逻辑。试驾行为说明用户已经走到线下环节了兴趣程度当然比单纯浏览高得多。如果你做的是其他行业同样要梳理自己业务的关键转化节点并合理分配权重。3.5 冷启动问题的处理ALS模型对老用户推荐效果好但新用户没有任何行为数据推荐结果会失效。我的方案是对新用户走热度推荐策略直接推荐当前销量最好、口碑评分最高的车型用户产生一定量的浏览行为后比如浏览了3辆车再切换到ALS个性化推荐。实现上SpringBoot里判断用户是否有历史行为记录没有就走热度推荐接口public ListCarModel getRecommendation(String userId) { if (redisTemplate.hasKey(rec:user: userId)) { // 直接用缓存中的推荐结果 return getFromRedis(userId); } // 检查用户是否有行为数据 int behaviorCount behaviorLogMapper.countByUserId(userId); if (behaviorCount 3) { // 冷启动返回热门车型 return carModelMapper.getHotModels(10); } // 正常推荐流程 ListCarModel recommendList getFromMySQL(userId); if (recommendList null || recommendList.isEmpty()) { return carModelMapper.getHotModels(10); } return recommendList; }4. 大数据分析模块实现4.1 销售趋势分析推荐系统只是其中一部分客户还要求做销售数据分析看板这个我同样用Spark实现。主要分析维度有按品牌、车型、价格区间的销量TOP排行按月/季度/年度的销量趋势对比不同城市、年龄段的购车偏好库存积压预警连续N天未出售的车型Spark分析任务每天凌晨2点启动批量处理前一天的全量数据结果写入MySQL的分区统计表。SpringBoot的报表接口只做查询展示逻辑非常简单但展示的数据背后都是Spark跑出来的聚合结果。4.2 热门车型画像用Spark SQL做聚合代码很简洁val hotCarAnalysis spark.sql( SELECT car_model_id, COUNT(DISTINCT user_id) as view_count, SUM(CASE WHEN behavior_type 4 THEN 1 ELSE 0 END) as inquiry_count, SUM(CASE WHEN behavior_type 5 THEN 1 ELSE 0 END) as test_drive_count, AVG(page_duration) as avg_duration FROM cleaned_behavior_log WHERE behavior_date 2024-06-01 GROUP BY car_model_id ORDER BY view_count DESC LIMIT 50 )这个结果可以用于管理后台的热门车型榜同时也可以为推荐系统冷启动提供热度评分数据。4.3 用户画像分析用户画像是大数据分析的另一个核心输出。我用Spark做用户分群通过用户的历史行为数据计算用户的品牌偏好、价位偏好、车辆类型偏好轿车/SUV/MPV等标签。具体的做法是把用户的各个偏好打分存成一个JSON格式的标签字段{ brand_pref: [大众, 丰田], price_range: [12, 18], car_type: SUV, purchase_intent: high }这个标签数据在推荐结果排序时非常有用。比如ALS推荐了10款车我再根据用户的价位偏好和车型偏好把同价位的SUV车型排序提到前面。相当于在协同过滤之上加了一层规则式过滤实测CTR提升了不少。5. SpringBoot与Spark的整合实践5.1 数据读取链路设计SpringBoot应用本身不直接参与Spark计算。我设计了两条数据读取链路链路一Spark计算后的结果表存在MySQL中SpringBoot通过MyBatis查询读取写入Redis缓存。这个链路用于推荐列表和报表展示。链路二SpringBoot接收实时行为日志异步发送到Kafka供Spark任务消费。这个链路用于日志采集与链路一互不影响。这样设计的好处是Spark的批处理计算和SpringBoot的在线服务完全解耦。即使Spark集群某一晚挂掉了SpringBoot服务依然能正常提供昨天的缓存数据不影响用户体验。等Spark恢复后重新跑一遍分析任务把结果更新回MySQL即可。5.2 SpringBoot整合KafkaSpringBoot发送行为日志到Kafka用异步线程池处理Configuration public class KafkaConfig { Bean public KafkaTemplateString, String kafkaTemplate() { return new KafkaTemplate(producerFactory()); } } Service public class BehaviorLogProducer { Autowired private KafkaTemplateString, String kafkaTemplate; Async(behaviorLogExecutor) public void sendBehaviorLog(BehaviorLog log) { String json JSON.toJSONString(log); kafkaTemplate.send(car_behavior_log, log.getUserId(), json); } }注意点Async必须配合线程池配置否则默认的SimpleAsyncTaskExecutor每次都会新建线程并发高的时候会创建上千个线程直接把内存打爆。踩过这个坑后我配置了一个核心线程数20、最大线程数100、队列容量1000的线程池效果稳定。5.3 Spark定时任务的调度Spark任务的调度我用的是SpringBoot的Scheduled注解配合分布式锁防止多节点重复执行Component public class SparkTaskScheduler { Scheduled(cron 0 0 2 * * ?) // 每天凌晨2点执行 public void runSparkTask() { // 尝试获取分布式锁 boolean locked redisTemplate.opsForValue().setIfAbsent(lock:spark:task, 1, 60, TimeUnit.MINUTES); if (!locked) { log.warn(Spark task already running on another node); return; } try { ProcessBuilder pb new ProcessBuilder(spark-submit, --class, com.car.spark.RecommendationJob, --master, yarn, /opt/spark-jobs/car-sales-spark.jar); Process process pb.start(); // 日志流处理... int exitCode process.waitFor(); log.info(Spark task finished with exit code: {}, exitCode); } finally { redisTemplate.delete(lock:spark:task); } } }用spark-submit方式提交任务比用SparkSession直接嵌入SpringBoot更合适。后者会让SparkContext常驻内存资源占用太高太浪费。5.4 Spark读取Redis实现推荐结果生成后如果把几十万用户的推荐列表全部写回MySQL每次查询性能会受影响。更合理的做法是把推荐列表写入Rediskey设计为rec:user:{userId}value是JSON数组。Spark写Redis有两种常见方式我分别测试过第一种是foreachPartition在RDD分区内批量写入recommendations.foreachPartition { partition val config new JedisPoolConfig() val pool new JedisPool(config, redis-host, 6379) val jedis pool.getResource partition.foreach { row val userId row.getAs[String](userId) val carList row.getAs[Seq[Row]](recommendations) val json carList.map(...).mkString([, ,, ]) jedis.set(srec:user:$userId, json) jedis.expire(srec:user:$userId, 86400) } jedis.close() pool.destroy() }第二种是用Redis的pipeline优化批量写入。实测下来10万条推荐结果foreachPartition逐条写需要55秒用pipeline可以压缩到8秒以内。强烈建议用pipeline。6. 系统部署与调优6.1 Spark集群搭建要点我会简单讲一下集群搭建的关键配置具体安装步骤系统环境不同会有差异重点是要理解哪些参数决定了系统性能。我用的集群是三台8核32G的服务器一台master、两台worker。核心配置有三个文件spark-defaults.conf、spark-env.sh、workers。spark-defaults.conf里我重点调了这几个参数spark.serializerorg.apache.spark.serializer.KryoSerializer spark.memory.offHeap.enabledtrue spark.memory.offHeap.size8G spark.sql.shuffle.partitions100Kryo序列化器比Java默认的序列化方式快10倍以上特别是数据量大时效果极其明显。offHeap内存可以绕过JVM GC的压力但要注意别把offHeap设太大否则影响其他进程。spark-env.sh里配置了JVM参数SPARK_WORKER_MEMORY8G SPARK_WORKER_CORES4 SPARK_EXECUTOR_MEMORY4G SPARK_EXECUTOR_CORES2 SPARK_DRIVER_MEMORY2G这里要说明executor内存不能超过worker内存减去系统开销否则会申请不到资源。比如worker内存8Gexecutor最多给6G比较安全。6.2 Spark作业参数调优同一份Spark代码参数配不好和配好性能差距可能达到数十倍。我拿ALS训练任务举例说明调优思路。最初我用默认参数跑30万用户的数据集跑了17分钟。后来做了三处调整缩短到3.5分钟第一设置合理的executor数量和每executor的核数。我的集群是三台机器每台分配2个executor每executor 2核。这样并行度是12对于ALS这种计算密集型的任务比较合适。spark-submit --class com.car.spark.ALSJob \ --master yarn \ --num-executors 6 \ --executor-memory 4G \ --executor-cores 2 \ /opt/spark-jobs/car-sales-als.jar第二用persist(StorageLevel.MEMORY_AND_DISK)缓存中间结果。ALS每次迭代都需要复用同一份特征矩阵如果每次都重新加载IO开销巨大。第三调整spark.sql.shuffle.partitions。这个参数决定shuffle时产生的分区数太小时task处理数据量过大太大会产生大量小文件。经验值是集群总核数的2到3倍。6.3 SpringBoot服务调优SpringBoot侧的调优重点在缓存和连接池。推荐接口的QPS压力主要在Redis读取所以我把缓存设计成两级本地Caffeine缓存 Redis缓存。本地缓存保存最近10分钟的热门推荐数据Redis保存全量推荐数据。这样热门用户请求直接在本地缓存命中Redis的负载降到了原来的30%。MySQL连接池用HikariCP推荐结果表加了联合索引(user_id, car_model_id)单次查询稳定在3ms左右。Kafka消费者线程数也要合理设置。我最初用的是默认配置一条一条处理消息吞吐量太低。后来调整为批量消费模式每次拉取500条一次性处理吞吐量提升了大约8倍。7. 常见问题与排查技巧7.1 Spark作业频繁OOM这是最常遇到的问题。解决方案不是简单调大内存而是先定位是driver端还是executor端OOM。如果错误日志里出现java.lang.OutOfMemoryError: Java heap space同时报错在executor需要考虑降低每个executor的内存占用合理吗大多数情况下是spark.sql.shuffle.partitions设置太小导致每个task处理的数据量过大。把该参数调大让数据分散到更多task里处理。如果是driver端OOM常见于collect()操作把大量数据拉回driver。能用save到文件或MySQL替代的就不要用collect()。7.2 推荐结果为空推荐列表为空有几种情况模型训练时设置了setColdStartStrategy(drop)新用户直接被丢弃了。这个时候SpringBoot要走兜底逻辑热度推荐训练数据太稀疏很多用户的行为数据不足。处理办法是调整评分阈值或者对行为数据做补充比如把浏览时长超过30秒的行为也算有效评分MySQL中推荐结果表没有数据。排查Spark任务是否成功执行查看exit code和日志7.3 Spark任务跑完但没更新数据这个问题我排查了很久。Spark任务日志显示success但MySQL里的推荐结果还是旧的。原因是Spark写MySQL时用了save modeJava进程里设置成了Append而数据源读取时没有去重导致重复执行多次后数据量翻倍了。解决方案写MySQL前先按userId去重或者清空当天的分区再写入。我现在用的方案是每天先执行DELETE FROM recommend_result WHERE biz_date ?然后再插入新数据保证幂等性。7.4 Kafka消费者重复消费在SpringBoot整合Kafka时如果消费者没有正确提交offset服务重启后会重复消费同一批消息。数据重复后会影响推荐结果计算的准确性。解决方案是把消费端设置为手动提交offsetKafkaListener(topics car_behavior_log, groupId car-sales-group) public void listen(ListConsumerRecordString, String records, Acknowledgment ack) { try { for (ConsumerRecordString, String record : records) { handleRecord(record); } ack.acknowledge(); // 等这一批处理完再提交offset } catch (Exception e) { // 处理失败不提交offset下次启动会重新消费 log.error(Kafka batch process error, e); } }这也就是典型的至少一次语义。如果对结果准确性要求极高需要配合幂等写入设计如果只是做分析不提交offset导致重复消费的问题通过清洗阶段的去重就能解决方案取舍看具体场景。7.5 问题速查表问题表现可能原因处理方法Spark任务OOMshuffle分区数过小调大spark.sql.shuffle.partitions推荐结果为空模型冷启动被丢弃兜底热度推荐推荐结果不准评分权重不合理调整行为评分映射接口响应慢缓存未生效检查Redis连接和缓存设置Spark任务重复跑没有分布式锁使用Redis分布式锁Kafka消费重复offset提交方式错误改为手动批处理提交8. 实际效果与后续优化方向8.1 系统上线后实测数据系统上线跑了一个月我整理了一些关键数据供大家参考推荐接口平均响应时间8.6msP99响应时间35ms。推荐位的点击率比原来的列表页提升了23%试驾转化率提升了11%。大数据分析看板的统计结果与线下财务报表的误差控制在2%以内基本可以替代人工Excel统计。Spark任务每天的执行时间也从刚开始的45分钟优化到了15分钟左右完全赶在早上8点业务高峰前完成数据更新。8.2 后续可以做的优化第一引入实时推荐。现在用户当天浏览的行为要到第二天凌晨才能反映在推荐结果里。如果想做到实时反馈可以接入Structured Streaming实时流计算结合用户当前的浏览动作动态更新推荐列表。但成本会高不少需要权衡。第二融合更多数据源。现在只用行为数据做推荐后续可以把车型参数、用户反馈评论、售后维修记录都纳入分析。如果用户的车型在反馈中提到质量问题推荐系统可以适当降低该车型的排序权重。第三做AB实验平台。推荐系统的价值要通过对比才能体现建议搭一套简单的AB分流框架对比不同算法版本、不同权重策略的效果让优化决策用数据说话而不是靠个人感觉调整算法参数。以我个人的体会来说这个项目最有价值的部分不是技术本身有多难而是把推荐算法和大数据分析这两套东西真正落地到了汽车销售的业务场景里。技术选型只是手段最终能帮销售团队把转化率提上来让管理层每天打开看板就能掌握全局这才是系统存在的意义。如果你们公司也有类似的业务场景可以参考这套方案从离线分析做起逐步迭代一定会看到效果。
上一篇/下一篇内容由系统自动关联
返回资讯列表 →