ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Hadoop+Spark+SpringBoot在线广告推荐系统实战全解析

Hadoop+Spark+SpringBoot在线广告推荐系统实战全解析 做广告流量变现相关项目的朋友应该都有体会推荐系统听起来高大上但真正落地时最折磨人的不是算法本身而是整条链路怎么打通。最近刚完成一个基于HadoopSparkSpringBoot的在线广告推荐系统项目从购药日志处理到ALS协同过滤模型训练再到SpringBoot接口服务和可视化大屏整条数据链路完整跑通。这篇文章我把整个项目从架构设计、环境搭建、算法实现到前后端联调的细节全部分享出来包括我踩过的坑和最终采用的解决方案。适合正在做大数据方向课程设计、毕业设计或者公司里准备搭建广告推荐系统的开发同学参考。1. 项目全景拆解在线广告推荐系统的业务逻辑与技术选型1.1 核心需求与数据流转路径在线广告推荐平台的业务场景很直接用户在浏览媒体内容时系统需要根据用户的历史行为把最有可能被点击的广告推送到用户面前。听起来简单但要支撑这个功能需要处理的数据链路相当长。这个项目里我确定的整体数据流转路径是埋点日志产生用户行为数据包括广告曝光、点击、浏览时长等日志文件落入HDFS存储形成数据底座Spark定时任务从HDFS读取日志清洗后用于ALS模型训练模型产出推荐结果写入MySQL供服务层调用SpringBoot后端读取MySQL中的推荐结果对外提供REST API可视化大屏通过后端API获取统计数据渲染运营监控看板这个流程里每个环节的输入输出都是明确的HDFS管存储Spark管计算MySQL管结果落地SpringBoot管服务化大屏管展示。各司其职互不干扰排查问题时也能快速定位。做课程设计或者毕业设计的同学我建议你也按照这个链路来设计。不要试图把Spark任务直接塞进SpringBoot里跑那样项目看起来是一个整体实际操作时会有很多麻烦后面我会详细解释原因。1.2 技术选型背后的权衡逻辑为什么选择HadoopSpark这套组合而不是更轻量级的方案这个问题的答案要回归到业务场景本身。广告日志是典型的海量小文件数据每天几百万条甚至上亿条都很正常。HDFS在处理这类海量文件存储上有天然优势扩展性也足够好后续数据量增长只需增加节点。Hadoop的生态配套也成熟不需要额外造轮子。Spark的选型则是看中了它在迭代式计算上的性能。ALS协同过滤算法需要反复在用户-物品矩阵上进行多轮迭代计算每一轮都要遍历全部数据。用传统的MapReduce做这种计算每轮迭代都涉及磁盘读写性能损失非常大。换用Spark的内存计算模型之后中间结果直接驻留在内存里同样是跑10轮迭代时间能缩短一个数量级。SpringBoot在这个体系里的定位是服务化接入层。广告推荐结果最终要提供给业务端调用需要用稳定的HTTP接口输出。SpringBoot生态成熟、上手快做REST API非常顺手。这里有一个关键的设计原则我在项目里始终坚持计算层和服务层严格分离。Spark任务跑完之后把结果沉淀到MySQLSpringBoot只负责从数据库读取数据两者之间没有直接调用关系。这样分开带来的好处是Spark集群抖动或者正在跑批的时候用户端的推荐接口依然能正常返回上一次的推荐结果线上服务的稳定性有保障。2. Hadoop与Spark环境搭建从伪分布式到集群的实战选择2.1 Hadoop伪分布式搭建的关键配置Hadoop环境搭建是很多初学者第一道坎。这个项目为了保证能在一台机器上完整验证流程我采用了伪分布式部署模式。所谓伪分布式就是所有Hadoop守护进程NameNode、DataNode、ResourceManager、NodeManager都跑在同一台中CPU上但配置方式和真正集群没有本质区别改一下IP配置就能横向扩展。这里我总结几个容易踩坑的配置点版本选择上我用的是Hadoop 3.3.4配合JDK 1.8。Hadoop 3.x对JDK有硬性版本要求如果JDK版本不对启动时会直接报错。另外新版本的Java再去运行老版本的Hadoop经常会遇到一些莫名其妙的反射调用异常直接用我这对组合最省心。免密登录这个步骤经常被跳过去但实际上非常重要。Hadoop进程之间需要通过SSH互相通信如果不配置免密启动过程中会反复报错Connection refused。配置方式很简单ssh-keygen -t rsa -P -f ~/.ssh/id_rsa cat ~/.ssh/id_rsa.pub ~/.ssh/authorized_keys chmod 600 ~/.ssh/authorized_keys三个核心配置文件的参数设置也分享一下。core-site.xml里要指定NameNode的地址和HDFS临时目录hdfs-site.xml里设置副本数和NameNode数据目录yarn-site.xml里配置资源管理器地址。伪分布式模式下副本数设置为1就够了保持默认的3会浪费磁盘空间。启动前有一件事千万别忘格式化NameNode。很多新手第一次启动失败都是因为改了配置文件但没有重新格式化。正确的操作顺序是改完配置之后执行hdfs namenode -format注意这个命令只需要执行一次。后期如果反复格式化会导致NameNode的namespace ID和DataNode不一致DataNode启动时会被拒绝注册。2.2 Spark On YARN部署与资源参数调优Spark的部署模式我选了YARN模式而不是Standalone。原因很简单项目里既有Spark任务又有MapReduce任务让YARN统一管理集群资源多个计算框架可以共享集群资源不会出现某一时间段Spark占满资源导致Hadoop任务饿死的情况。Spark接入YARN需要在spark-env.sh里配置两个关键项export JAVA_HOME/usr/local/java/jdk1.8.0_202 export HADOOP_CONF_DIR/usr/local/hadoop/etc/hadoop这里有个细节HADOOP_CONF_DIR一定要配否则Spark任务在YARN模式下找不到HDFS的配置信息连接HDFS时会报错。验证环境是否打通可以先启动一个Spark Shell试试spark-shell --master yarn --deploy-mode client能进入Scala交互式环境说明Spark和Hadoop的整合已经没问题了。如果在这个环节卡住最常见的两个原因是一是YARN的ResourceManager没启动检查一下jps进程列表里有没有ResourceManager二是内存参数配置不对导致Spark任务在申请不到Container就被kill了。伪分布式环境下跑Spark任务内存配置的坑最典型。默认情况下YARN的NodeManager可分配内存有限而Spark Executor动辄申请几个GB的堆内存抢不到Container就直接报错。我调整了yarn-site.xml里的几个参数property nameyarn.nodemanager.resource.memory-mb/name value8192/value /property property nameyarn.scheduler.maximum-allocation-mb/name value8192/value /property这组配置的意思是NodeManager节点上可供YARN分配的物理内存总额为8GB单个任务最大可申请8GB。这样Spark任务申请Executor内存时就不会因为配额不足而被拒绝。3. 推荐算法的核心Spark ALS协同过滤的原理与落地3.1 从用户行为到评分矩阵ALS算法全称Alternating Least Squares交替最小二乘法是协同过滤领域最经典的矩阵分解模型。要理解这个算法在广告推荐中怎么用要先搞清楚输入数据是什么。推荐场景里算法需要的是一个用户对物品的评分矩阵。在电影推荐里评分就是用户给电影打的1到5分但在广告推荐场景用户不会主动给广告打分所以要从行为数据里构造出评分。我采用的是将点击次数和浏览时长加权合成综合评分score 0.7 * click_count 0.3 * normalized_view_time为什么用综合评分而不是直接用点击次数因为单看点击次数有偏差有些用户只看广告但不点击浏览时长越长说明广告内容越符合他的兴趣反过来有人误点了广告但立刻关掉这种行为的价值也有限。加权重化之后评分对用户兴趣的刻画更加平滑。数据成形后ALS的核心是将稀疏的评分矩阵分解成两个低秩矩阵的乘积用户因子矩阵和广告因子矩阵。两个矩阵在k维特征空间内逼近原始评分矩阵。这样做的好处是即便用户和广告之间没有直接交互记录只要它们在因子空间上距离近就能推断出用户对广告的潜在偏好。3.2 训练集与测试集划分的核心细节数据预处理和划分阶段我遇到一个非常经典的信息泄漏问题。最初我直接用随机切分的方式把用户行为数据分成训练集和测试集结果模型评估的RMSE非常漂亮但上线后推荐效果一塌糊涂。后来排查原因发现同一个用户的记录同时出现在训练集和测试集里模型在训练时已经见过这个用户的偏好测试时自然表现得很好。解决办法是按用户维度进行切分确保同一个用户的数据只出现在一侧// 按用户分组后将每个用户的记录按比例拆分为训练集和验证集 val splits data.randomSplit(Array(0.8, 0.2), seed 42L)这里还有个细节点随机种子要保持固定否则每次运行结果不一致不方便复现实验效果。为了让代码有更好的可复现性我在所有涉及随机操作的环节都设置了固定的seed值。3.3 ALS训练与推荐生成的代码落地模型训练部分直接使用Spark MLLib提供的ALS算子。参数的选择对整个效果影响很大我用过几个不同参数组合做对比import org.apache.spark.ml.evaluation.RegressionEvaluator import org.apache.spark.ml.recommendation.ALS val als new ALS() .setMaxIter(10) // 最大迭代轮数 .setRegParam(0.1) // 正则化系数 .setRank(10) // 因子矩阵维度 .setUserCol(userId) .setItemCol(itemId) .setRatingCol(rating) .setColdStartStrategy(drop)这个参数组合里正则化系数0.1用于防止过拟合rank10意味着用户和广告都在10维空间里做特征表达。rank太小时模型的表达能力不足预测值会趋于一致rank太大则计算量大而且容易过拟合训练数据。项目开始时可以先用小rank跑通再逐步调大对比效果。评估指标我用的是RMSE均方根误差。公式为val predictions model.transform(test) val evaluator new RegressionEvaluator() .setMetricName(rmse) .setLabelCol(rating) .setPredictionCol(prediction) val rmse evaluator.evaluate(predictions)RMSE的意义在于描述预测评分和真实评分之间的偏移程度。数值越小说明模型的预测越精准。推荐结果的生成我用了recommendForAllUsers方法给每位用户生成10条广告推荐。这个操作会基于所有用户的因子矩阵和所有物品的因子矩阵做矩阵运算输出用户最感兴趣的TopN物品列表val recommendations model.recommendForAllUsers(10)产出的结果会写入MySQL表字段包括userId、itemId、预测评分、推荐生成时间。这步操作本质上就是把计算层的结果落到存储层供下游服务读取。写库时要注意批次大小控制一次性写太多容易把数据库连接池打满推荐用foreachBatch分批写入。4. SpringBoot后端服务设计推荐结果如何变成可用的接口4.1 项目分层结构与核心接口定义SpringBoot后端的职责很清晰对外提供推荐接口和运营统计接口不负责跑任何计算任务。项目结构采用标准的Mapper-Service-Controller三层架构com.example.adrec ├── controller │ ├── RecommendController.java │ └── StatsController.java ├── service │ ├── RecommendService.java │ └── StatsService.java ├── mapper │ ├── RecommendMapper.java │ └── StatsMapper.java ├── entity │ ├── RecommendResult.java │ └── AdStats.java └── config └── WebConfig.java接口定义上我推荐用RESTful风格。项目里设计了一组对外APIGET /api/recommend/{userId}按用户ID返回广告推荐列表GET /api/stats/overview返回大屏所需的汇总指标GET /api/stats/trend返回分时段点击趋势数据这样划分的好处是推荐场景和数据展示场景分离各自有独立的接口维护后续的迭代升级。推荐接口返回的JSON结构也做了一致化处理统一包含code、message、data三层结构。data里是推荐列表每个广告对象包含广告位ID、广告创意ID、推荐分数等字段。4.2 计算结果的存储选型与读取策略关于Spark计算结果落地存储选型我对比过MySQL和Redis两个方案。最终选了MySQL核心理由是推荐结果一天只更新一次频率很低而且MySQL方便Windows浏览器直接查看数据排查问题的时候SQL查询非常灵活。但这里必须处理好一个冷启动的问题。新注册的用户没有历史行为数据ALS模型给不出推荐结果。如果不做处理接口返回空列表前端拿到空数据整个推荐位就空了。我在Service层加了一个兜底策略if (recommendList null || recommendList.isEmpty()) { // 返回全局热门广告TopN return hotAdService.getGlobalHotAds(10); }热门广告的计算同样是Spark任务完成的按每个广告的曝光到点击的转化率排序取TopN。这样冷启动用户虽然拿不到个性化推荐但至少能看到平台当前最热的广告内容不至于页面空白。另一个要注意的点是MySQL查询性能。推荐结果表随着时间增加会快速膨胀如果不在查询字段上建立索引接口响应时间会越来越长。我给userId和recommend_time字段建了联合索引单次查询耗时从几百毫秒降到了几十毫秒级别ALTER TABLE recommend_result ADD INDEX idx_user_time (user_id, recommend_time);4.3 接口联调中的性能与兼容性处理联调阶段遇到3个问题值得大家留意。JSON字段命名不一致问题。Spark侧写出的字段名是user_idJava实体类里用的是userId。两边对不上解析的时候会出现字段丢失前端拿到的都是null。我统一用Jackson注解在实体类上做映射JsonProperty(user_id) private Long userId;接口响应慢的问题。排查下来发现是数据库查询时间过长而且同一时刻大量用户请求击穿了MySQL的连接池。解决方案是引入Spring Cache做了一层缓存把推荐结果缓存5分钟缓存命中的请求直接返回不查数据库Cacheable(value recommendCache, key #userId) public ListRecommendResult getRecommendList(Long userId) { ... }浏览器跨域的问题。后端没有做任何跨域处理时前端调用接口会被CORS策略拦截。在WebConfig里加一个全局CORS配置即可Configuration public class WebConfig implements WebMvcConfigurer { Override public void addCorsMappings(CorsRegistry registry) { registry.addMapping(/api/**) .allowedOrigins(*) .allowedMethods(GET, POST, PUT, DELETE); } }5. 可视化大屏实战从数据到图表的完整链路5.1 大屏指标设计与数据聚合策略可视化大屏不是简单地把几个图表拼在一起而是要回答运营最关心的问题。我在这个项目里把大屏拆成了几个核心模块每个模块回答一类问题。顶部核心指标区展示总曝光量、总点击量、整体CTR和平均转化率这些是广告投放效果的一级指标运营每天早上第一眼看的就是这组数字。中间区域放24小时点击趋势折线图反映用户的活跃时段分布运营可以据此调整广告的排期策略。右侧放广告位点击排行榜和广告主消耗榜帮助商务团队评估不同广告位的商业价值。底部是地域分布图展示不同区域的广告点击热度。为了支持这些大屏展示Spark批处理作业在训练模型的同时会把聚合统计数据计算出来写入MySQL中一张独立的聚合表。为什么要单独建表而不是直接查明细因为明细数据可能上千万条每次大屏加载都全表扫描数据库压力太大页面也卡得没法看。聚合表的数据量非常小是Day维度已经汇总好的值接口查询和传输都很快。5.2 ECharts大屏布局与动态刷新前端实现我选择了Vue配合ECharts这两者的组合灵活且社区生态好。大屏的自适应是常见问题我采用了scale缩放方案固定设计稿尺寸1920x1080运行时动态计算屏幕实际尺寸与设计稿的比例对整个大屏容器做transform缩放function scaleScreen() { const ratioX window.innerWidth / 1920; const ratioY window.innerHeight / 1080; const scale Math.min(ratioX, ratioY); document.getElementById(screen).style.transform scale(${scale}); } window.addEventListener(resize, scaleScreen);这种方案能保证大屏在不同分辨率屏幕上都以完整比例显示不会出现错位或溢出。数据刷新部分大屏通过ak定时轮询后端API每30秒刷新一次。前端先拉取数据再用ECharts的setOption方法更新各大图表。大屏配色方面推荐深色底配亮色图表的方案。我用的配色是背景#0f1130主色#2f6bff辅助色#f5a623和#00d4c8。深色背景能有效减少大屏在弱光环境下的刺眼感亮色图表数据更容易看清。6. 项目调试与排错实录环境、数据与联调问题的完整排查6.1 Hadoop与Spark环境类问题排查整个项目开发过程中环境类问题占据了调试时间的大头。我把两个最常见的坑整理如下大家遇到可以直接对照操作。第一个是NameNode启动失败。表现形式是执行start-dfs.sh之后jps命令看不到NameNode进程日志里报Incompatible namespaceIDs错误。这个问题的根源在于格式化操作被多次执行NameNode的namespace ID和DataNode的不一致。解决方案是彻底清理数据目录后重新格式化# 1. 先停掉所有Hadoop进程 stop-all.sh # 2. 删除数据目录注意路径要与hdfs-site.xml配置一致 rm -rf /usr/local/hadoop/data/ # 3. 重新格式化 hdfs namenode -format # 4. 再次启动 start-all.sh第二个坑是Spark任务在YARN模式下被频繁kill日志里显示Container killed on request. Exit code is 143。检查了YARN日志后确认是物理内存超限。伪分布式机器只有16GB内存YARN默认参数给Executor分配的内存太少任务执行时超过限制就被杀掉。把yarn-site.xml的内存参数调高之后问题解决。6.2 数据与模型参数调优排查数据类问题里面ALS训练报错的情况最多。一个典型的场景是rating列包含空值ALS要求评分不能为空遇到空值直接抛异常。解决方案是在训练前过滤val cleanData rawData.filter(col(rating).isNotNull col(userId).isNotNull col(itemId).isNotNull)另一个常见是预测结果全部为NaN。这个问题往往是正则化系数处理不当导致梯度爆炸或参数发散。regParam建议初始设为0.1如果出现发散再逐步调小到0.01。推荐效果不好所有用户拿到的推荐结果高度相似基本可以判断是rank参数太小因子空间维度不足。我在实验中把这从5调到10之后推荐结果的区分度明显提升。6.3 前后端联调与部署上线经验联调阶段最容易被忽略的是数据库连接池配置。SpringBoot默认的连接池参数在并发请求上来后会经常报连接获取超时。我给HikariCP做了调整设置最大连接数为50连接超时时间为30秒设置之后接口的稳定性大幅提升。上线部署时也提醒一句确保Spark任务执行的目录和日志目录有足够磁盘空间我遇到过跑批任务时数据量大把磁盘塞满导致写入MySQL失败的情况。给Spark作业日志设置滚动保留策略或者定期清理历史日志文件能避免这类问题。还有就是在Spark任务里最好加上重试机制。HDFS在任务执行期间可能出现短暂的节点抖动一次性任务失败后手动重跑的成本很高。给关键任务加上失败自动重试1到2次的配置可以大幅提升容错性。做这个基于HadoopSparkSpringBoot的在线广告推荐系统最深的感触是整个项目的价值在于把各组件真正串成了一个完整闭环。数据从日志文件一步步变成推荐结果和大屏图表每一条数据链路里都藏着环境配置、格式转换、性能调优等大量细节。别指望一次跑通先让每个环节单独验证再打通链路是最高效的推进方式。最后分享一个小技巧每完成一个环节就把计算结果落成CSV用肉眼检查几行确认数据形态和预期一致之后再进入下一个环节。这套项目做完之后我后续的扩展方向是引入Kafka做实时数据接入配合Spark Streaming给用户做秒级推荐刷新再叠加用户画像特征做混合推荐。有了现在这套稳定的离线链路打底这些升级都只是时间问题。
返回列表