news 2026/9/30 15:07:06

基于Spark的在线广告推荐系统实战:从ETL到可视化大屏

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
基于Spark的在线广告推荐系统实战:从ETL到可视化大屏

说句实话,第一次看到“基于 Spark 的在线广告推荐系统”这个项目名时,我也觉得它挺唬人。在线广告推荐,听着像大厂算法团队才能碰的东西;但把它落到 Hadoop + Spark + Spring Boot 这套技术栈上,它其实就是一个特别典型的“大数据离线计算 + 业务接口 + 可视化展示”的组合题。我按源码、文档、调试、可视化大屏这四个维度把整个项目完整跑通过一遍,这里把真实的搭建过程、代码组织逻辑和排查经验一次性写清楚。无论你是做课程设计、毕业设计,还是想补一套完整的大数据项目经历,这篇文章都能当一份操作手册直接参考。

1. 项目整体设计与技术选型

1.1 技术栈组合的逻辑

先回答一个问题:为什么偏偏是 Hadoop、Spark、Spring Boot 这三个东西凑在一起?

在线广告推荐系统最核心的需求其实只有三个:存得下、算得快、发得出。用户的点击曝光日志是海量的,需要一套分布式存储把它接住,这是 Hadoop 里 HDFS 的活儿;原始日志进来之后要清洗、聚合、训练模型、生成推荐列表,这是 Spark 的活儿;推荐结果不能直接丢给用户,需要一个后端服务把结果封装成接口,这是 Spring Boot 的活儿。可视化大屏则是把最终数据用图表展示出来,解决“做完不知道效果如何”的问题。

用生活化的类比来说:Hadoop 是仓库,Spark 是加工厂,Spring Boot 是客服窗口,可视化大屏是车间外的仪表盘。仓库收货,工厂加工,客服把加工好的东西送出去,仪表盘让人看清整条产线是否正常。四个角色各管一段,边界清晰。这也是我最终没有把 Spark 作业嵌进 Spring Boot 进程的原因——API 服务要保持轻量,计算任务交给 YARN 调度,两者一分离,部署和排查都轻松很多。

如果你只是在单机上做课程设计,Hadoop 做成伪分布式就够用;如果你要撑起更大数据量,再考虑把 Spark 改成 standalone 模式或者集群模式。记住一个原则:能用伪分布式跑通,就没有必要一上来就搞五台机器,环境复杂度往往是项目烂尾的头号原因。

1.2 在线广告推荐的完整链路

一个可复现的在线广告推荐系统,从数据产生到最终展示,完整链路大致是下面这七步:

  1. 前端投放页面或 SDK 埋点,上报广告曝光、点击、浏览时长等行为日志。
  2. 日志通过 Nginx 或 Flume 落到 HDFS 的指定目录,按日期分区存放。
  3. Spark 定时任务读取当日和历史 HDFS 日志,完成数据清洗、字段规整、用户画像聚合。
  4. 在清洗后的数据上训练或更新推荐模型,生成每个用户的广告 TopN 列表。
  5. 把推荐结果写入 Redis 或 MySQL,方便后端服务快速读取。
  6. Spring Boot 对外提供推荐接口和大屏统计接口,前端发起请求获取数据。
  7. 可视化大屏定时刷新,展示曝光量、点击率、消耗金额、热门广告排行等指标。

很多初学者容易把注意力全放在“算法”上,但实际项目中,链路通畅度比模型复杂度重要得多。一个基于统计规则或协同过滤的简单模型,只要能端到端跑通,就已经超过八成停留在 Demo 阶段的项目了。

1.3 源码目录怎么组织

源码交付最忌讳的就是所有文件堆在一个目录里,连哪个模块是干什么的都要靠猜。我在这个项目里用的是分层目录结构,看起来繁琐,但后期调试和写文档非常省事:

ad-recommend/ ├── data/ # 模拟日志和初始化数据 │ ├── ad_log_2025-01-01.json │ └── sql/ │ ├── admeta.sql # 广告主、广告位、广告物料表 │ └── dashboard.sql # 大屏统计结果表 ├── etl-spark/ # Spark 离线计算模块 │ ├── src/main/python/ │ │ ├── ad_etl.py # 日志清洗 │ │ ├── als_train.py # 协同过滤召回模型 │ │ └── result_writer.py # 写 Redis/MySQL │ └── pom.xml ├── ad-api/ # Spring Boot 服务模块 │ ├── src/main/java/com/lan/ad/ │ │ ├── controller/ # 推荐、大屏、用户接口 │ │ ├── service/ # 业务逻辑和缓存逻辑 │ │ ├── mapper/ # MyBatis-Plus Mapper │ │ ├── model/ # 实体和 DTO │ │ └── config/ # Redis、CORS、线程池配置 │ └── src/main/resources/ │ ├── application.yml │ └── mapper/ ├── ad-dashboard/ # 可视化大屏前端 │ ├── src/ │ │ ├── api/ # axios 封装 │ │ ├── components/ # KPI 卡片、图表组件 │ │ ├── views/Dashboard.vue │ │ └── utils/request.js │ └── package.json ├── docs/ │ ├── 部署文档.md │ ├── 接口文档.md │ ├── 数据字典.md │ └── 答辩问题整理.md ├── scripts/ │ ├── start-all.sh │ └── stop-all.sh └── README.md

这种结构让我在写文档时几乎不用额外思考,什么地方该写什么,目录里已经给了提示。尤其是 scripts 目录单独放出来之后,每次重新搭建环境只需要照着脚本顺序执行,省去了手动敲一堆命令的时间。

1.4 依赖版本的坑,提前说明白

大数据项目最怕版本漂移。我见过很多人把 Hadoop 2.x、Spark 3.x、Spring Boot 3.x 混在一起,最后全是 javax 和 jakarta 命名空间冲突、JDK 版本不兼容这类问题,项目还没开始写就已经劝退。这个项目我最终锁定的版本组合如下:

组件推荐版本说明
JDK1.8 或 11和 Spark 2/3 生态兼容性最好
Hadoop3.3.6内置 NameNode HA 支持,稳定
Spark3.3.4支持 PySpark 和 Scala,读 JSON 方便
Spring Boot2.7.18避免 3.x 的 Jakarta 迁移问题
MySQL8.0存广告元数据和统计结果
Redis7.x缓存推荐结果和用户实时标签

如果你非要用 Spring Boot 3.x,也不是不行,但需要特别注意两点:一是 MyBatis-Plus、Spring Data Redis 这些库要选适配 Jakarta 的版本,二是 Spark 客户端依赖如果打不进 Spring Boot 进程,就以独立任务方式运行。我实际做下来,Spring Boot 2.7.18 + JDK 8 是最省心的组合,没有之一。

2. Spark 数据清洗与推荐模型落地

2.1 广告点击日志字段设计

推荐模型的效果很大程度取决于日志质量。如果日志字段前后不一致,清洗脚本就是灾难。我在项目里把日志设计成 JSON 格式,每个对象代表一条曝光或点击记录,核心字段如下:

字段示例值含义
userIdU10001用户唯一标识
adIdA30021广告唯一标识
advertiserIdADV88广告主标识
categoryIdC_IT广告所属类目
channelId101投放渠道
deviceTypeandroid设备类型
province广东地域
isClick1是否点击,0/1
showTime2025-01-01 10:23:45曝光时间
clickTime2025-01-01 10:24:02点击时间
staySeconds18落地页停留时长
adTitle限时抢购广告标题
tokenxxxxxx埋点去重标识

日志格式一旦定下来,就尽量别改。你可以加字段,但不要调整已有字段的顺序或类型,因为下游 Spark 任务会严格按照这个 schema 做解析。我建议起步阶段用spark.read.json直接读,后面数据量大了再考虑用 Avro 或 Parquet。

2.2 Spark ETL 清洗代码

清洗是整个推荐质量的地基。我在项目里用 PySpark 写了一个 ETL 作业,大致逻辑如下:

from pyspark.sql import SparkSession from pyspark.sql import functions as F spark = SparkSession.builder \ .appName("ad_etl_clean") \ .master("yarn") \ .config("spark.sql.shuffle.partitions", "200") \ .getOrCreate() # 读取 HDFS 上的当天日志 raw_df = spark.read.json("hdfs://localhost:8020/user/ad_log/2025-01-01/*.json") # 1. 过滤明显异常数据 cleaned_df = raw_df.filter( F.col("userId").isNotNull() & F.col("adId").isNotNull() & F.col("isClick").isin([0, 1]) ) # 2. 去重,同一用户在 5 秒内对同一条广告的重复曝光只保留一条 cleaned_df = cleaned_df.dropDuplicates(["userId", "adId", "showTime"]) # 3. 统一时间格式,按小时提取分桶字段 cleaned_df = cleaned_df.withColumn( "showHour", F.date_format("showTime", "yyyy-MM-dd HH") ).withColumn( "clickFlag", F.col("isClick") ) # 4. 写出到 Parquet,按日期分区 cleaned_df.write \ .mode("overwrite") \ .partitionBy("showHour") \ .parquet("hdfs://localhost:8020/user/ad_cleaned/2025-01-01")

清洗逻辑看起来简单,但每一条都很关键。过滤isClick不在 0/1 范围的脏数据,是为了避免模型学到错误标签;去重是为了防止埋点重复上报导致指标虚高;按小时分区则是为后续大屏按时间维度统计做好准备。

如果你要直接跑,可以把master("yarn")改成master("local[*]"),先本地验证逻辑,再提交到集群。PySpark 的好处是本地和集群只需要换一行配置,语法可以完全一致。

2.3 基于 ALS 的召回模型

在线广告推荐不需要一开始就上深度学习。对课程设计或中小数据量场景,Spark MLlib 里的 ALS 协同过滤是最容易落地且效果稳定的方案。它的思路很简单:根据用户对广告的点击行为,把用户和广告映射到同一个隐语义空间,再通过矩阵分解找到“和这个用户历史行为相似的其他广告”。

ALS 默认适用于显式评分,但广告场景里我们只有“点了”和“没点”,属于隐式反馈。所以我把点击次数、停留时长构造成一个隐式评分,再用implicitPrefs=True训练:

from pyspark.sql import functions as F from pyspark.ml.recommendation import ALS # 构造用户 - 广告评分表 rating_df = cleaned_df.groupBy("userId", "adId").agg( F.sum(F.col("isClick")).alias("clickCnt"), F.avg(F.col("staySeconds")).alias("avgStay") ).withColumn( "score", F.when(F.col("clickCnt") > 0, 1.0 + F.log1p(F.col("avgStay"))).otherwise(0.0) ) # 划训练集和验证集 train, test = rating_df.randomSplit([0.8, 0.2], seed=42) # ALS 模型 als = ALS( userCol="userId", itemCol="adId", ratingCol="score", implicitPrefs=True, alpha=0.5, rank=10, maxIter=10, regParam=0.05, coldStartStrategy="drop" ) model = als.fit(train) # 为每个用户生成 10 条广告推荐 user_recs = model.recommendForAllUsers(10) user_recs.show(5, truncate=False)

coldStartStrategy="drop"一定要加。否则验证集里出现训练集没见过的用户时,Spark 会直接返回空值,后续写 Redis 时很容易抛异常。参数方面,rank控制隐向量维度,太小欠拟合,太大容易过拟合且吃内存;implicitPrefs=True表示我们输入的不是真实评分,而是隐式反馈;alpha控制正负样本权重,通常在 0.2 到 0.8 之间调整。

如果你的广告标题里有类目和关键词信息,还想做得再细一点,可以在 Spring Boot 里集成 HanLP 分词,对广告标题抽取关键词,做成内容特征做二次召回。这一块属于可扩展项,先不放进核心链路。

2.4 把推荐结果写回 Redis 和 MySQL

Spark 算出来的user_recs是一张 DataFrame,本质上是分布式的,不能直接让 Spring Boot 像查数据库一样查它。所以一定要把结果落到适合在线读的存储里。我采用双写方案:推荐列表写 Redis,统计指标写 MySQL。

写 Redis 时要注意分布式的连接池问题,不能每个分区都创建一个连接,最好用foreachPartition,让每个 Executor 里的分区统一写一批数据:

import json import redis def write_batch_to_redis(partition): r = redis.Redis(host="node01", port=6379, db=0) pipeline = r.pipeline() for row in partition: user_id = row["userId"] rec_list = [ {"adId": item["adId"], "score": round(item["rating"], 4)} for item in row["recommendations"] ] key = f"rec:user:{user_id}" pipeline.setex(key, 600, json.dumps(rec_list, ensure_ascii=False)) pipeline.execute() user_recs.foreachPartition(write_batch_to_redis)

TTL 设置为 600 秒,也就是 10 分钟过期一次。这样模型结果不会一直不变,即便 Spark 任务今天跑失败,推荐接口也还能用旧缓存顶上,不会直接返回空列表。

MySQL 里可以建一张recommend_result表,字段包含user_id、ad_id、score、update_time,作为 Redis 不可用时的兜底数据源。如果后面要上实时推荐,还可以把用户实时行为写入 Kafka,再用 Spark Structured Streaming 消费,这里先不做展开。

2.5 Hadoop 在项目里到底承担了什么

很多新手把 Hadoop 理解成一个单独的数据库,其实它是一整套基础设施。在这个项目里,它主要承担三件事:

第一是存储。清洗前后的广告日志全部放在 HDFS 里,相当于把文件系统的容量横向扩到几十上百台机器。第二是资源调度。Spark 作业跑在 YARN 上,YARN 分配 CPU 和内存给各个 Executor,这样多个计算任务才能共享同一批机器。第三是生态联动。Spark 原生支持读取 HDFS 路径,spark.read.json("hdfs://...")就是最直观的体现。

如果你在本地做伪分布式,Hadoop 的核心配置其实就三处:

<!-- core-site.xml --> <configuration> <property> <name>fs.defaultFS</name> <value>hdfs://localhost:8020</value> </property> </configuration>
<!-- hdfs-site.xml --> <configuration> <property> <name>dfs.replication</name> <value>1</value> </property> </configuration>

dfs.replication设成 1,是因为单机伪分布式只需要一份副本,设成默认的 3 反而会在写入时报错。

如果你要组建真正的多机集群,那就需要配置 NameNode 和 YARN ResourceManager 的地址,并且让所有机器的 SSH 互通。只有在做 NameNode 高可用 HA 的时候才需要引入 Zookeeper,课程设计阶段完全可以跳过。很多同学一上来就配 Zookeeper,结果三个节点互选主选不出来,最后连启动都卡住。我的建议是:单机伪分布式跑通,再考虑集群扩展,别一步登天。

3. Spring Boot 推荐服务接口实现

3.1 工程骨架和依赖

Spring Boot 在这个项目里不负责计算,它只做三件事:向外暴露推荐接口、提供大屏统计接口、管理广告元数据。因此它的 Maven 依赖也不需要引入 Spark 全套,那会让 jar 包体积暴涨,还容易产生冲突。

我建了一个标准 Web 工程,核心依赖如下:

<parent> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-parent</artifactId> <version>2.7.18</version> </parent> <dependencies> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-web</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-data-redis</artifactId> </dependency> <dependency> <groupId>com.baomidou</groupId> <artifactId>mybatis-plus-boot-starter</artifactId> <version>3.5.3</version> </dependency> <dependency> <groupId>mysql</groupId> <artifactId>mysql-connector-java</artifactId> <version>8.0.33</version> </dependency> </dependencies>

之所以不把 Spark 和 Hadoop 的客户端依赖放进 Spring Boot,是因为两个进程的生命周期完全不一样。Spark 作业是批任务,跑完就退出;Spring Boot 是常驻服务,需要稳定运行不被 OOM 干扰。把两者硬拆开,是项目稳定的第一步。

3.2 推荐接口怎么设计

推荐接口的路径我设计成/api/recommend/{userId},返回当前用户最感兴趣的 10 条广告。处理逻辑非常简单:

@RestController @RequestMapping("/api/recommend") public class RecommendController { @Autowired private StringRedisTemplate redisTemplate; @Autowired private AdMapper adMapper; @Autowired private RecommendService recommendService; @GetMapping("/{userId}") public Result<List<AdItem>> recommend(@PathVariable Long userId) { return Result.ok(recommendService.getRecommendAds(userId)); } }

核心逻辑在 Service 层。先查 Redis,命中就直接返回;没命中就查 MySQL 里的兜底广告列表。这个策略能保证即使 Spark 任务今天没跑,接口也不会挂:

public List<AdItem> getRecommendAds(Long userId) { String key = "rec:user:" + userId; try { String json = redisTemplate.opsForValue().get(key); if (StringUtils.hasText(json)) { return JSON.parseArray(json, AdItem.class); } } catch (Exception e) { log.warn("Redis 读取失败,走 MySQL 兜底", e); } return adMapper.selectRandomAds(10); }

这里最容易被忽略的一点是 Redis 的 value 序列化。如果你用默认的 JdkSerializationRedisSerializer,写入的是二进制对象,Spring Boot 自己读没问题,但大屏前端直接拿字符串就会看到一堆转义符。我在配置里统一改成 Jackson 序列化,key 用 String,value 用 String,这样所有客户端读到的都是纯 JSON 字符串。

3.3 大屏数据接口和跨域配置

大屏需要的是汇总结果,不是个性化推荐。所以我单独设计了一个DashboardController:

@RestController @RequestMapping("/api/dashboard") public class DashboardController { @Autowired private DashboardService dashboardService; @GetMapping("/overview") public Result<OverviewVO> overview() { return Result.ok(dashboardService.getOverview()); } @GetMapping("/trend") public Result<List<TrendItem>> trend() { return Result.ok(dashboardService.getHourlyTrend()); } }

大屏前端和后端端口不同,必须处理跨域。Spring Boot 2.7 里可以这样配置:

@Configuration public class CorsConfig implements WebMvcConfigurer { @Override public void addCorsMappings(CorsRegistry registry) { registry.addMapping("/api/**") .allowedOriginPatterns("*") .allowedMethods("GET", "POST", "PUT", "DELETE", "OPTIONS") .allowedHeaders("*") .allowCredentials(true) .maxAge(3600); } }

如果你的项目升级到 Spring Boot 3.x,记得把javax.servlet的写法改成jakarta.servlet,这是最容易踩的版本坑。

3.4 服务打包和部署

Spring Boot 工程用maven clean package打成可执行 jar,然后一条命令启动:

nohup java -Xms512m -Xmx1024m -jar ad-api.jar > api.log 2>&1 &

-Xms和-Xmx不建议设成一样,虽然防止了 JVM 动态扩容,但会让整个服务占用的内存固定在高位,和 Spark 作业抢资源。设成512m起、1024m封顶是一个比较稳妥的选择。

有人拿到别人打好的 jar 之后,会尝试用反编译工具去看内部类结构。我只能说,临时看一下方法的实现思路可以,但真正要改功能,还是老老实实拿到源码工程重新构建,否则改完的 class 无法维护,也没有配套的测试。这也是我在交付时坚持把源码结构做干净的原因。

4. 可视化大屏实现

4.1 大屏指标的口径要统一

大屏本身不产生数据,它只是把 MySQL 或者 Redis 里的统计结果画成图表。所以第一步不是写代码,而是把指标口径定下来。我项目里使用的核心指标如下:

指标计算公式说明
曝光数当日广告曝光日志总数按 token 去重
点击数当日 isClick=1 的日志总数点击埋点去重
点击率 CTR点击数 / 曝光数 x 100%保留两位小数
消耗金额按广告主点击单价计算单价来自广告元数据表
eCPM消耗金额 / 曝光数 x 1000千次曝光收益
活跃用户数当日有曝光或点击行为的用户数按 userId 去重

为什么单独说“口径”?因为我在调试时遇到过一个问题:Spark 清洗任务里的去重规则和 MySQL 统计 SQL 里的去重规则不一样,导致大屏上的点击率和 Excel 里手动算出来的结果对不上,最后查了半天才发现是埋点 token 去重逻辑没有同步。所以,所有指标的计算逻辑必须集中在一个地方,要么全部在 Spark 作业中算好,要么全部在 MySQL 里写死,前端只是拿数展示。

4.2 大屏页面布局和技术选型

前端我用了 Vue3 + Vite + ECharts,这是目前最主流也最容易上手的大屏组合。页面布局分四个区域:

  • 顶部:4 个 KPI 卡片,展示曝光数、点击数、点击率、消耗金额。
  • 左侧:每小时点击趋势折线图。
  • 中间:热门广告 Top10 柱状图。
  • 右侧:渠道占比环形图。

大屏适配是一个容易忽略的细节。我采用两层方案:第一层用rem单位布局,第二层用transform: scale做整体缩放,保证在不同分辨率的投屏上不会错位。如果你只是本地展示,可以直接定死一个 1920 x 1080 的尺寸,然后按比例缩放容器。

const baseFontSize = 16 document.documentElement.style.fontSize = baseFontSize + 'px'

4.3 ECharts 核心图表代码

下面这段是“热门广告 Top10”柱状图的完整逻辑。大屏的数据统一走 axios 请求后端接口,返回数组之后灌进 ECharts 的option:

import * as echarts from 'echarts' import { getTopAds } from '@/api/dashboard' export function initTopAdsChart(domId) { const chart = echarts.init(document.getElementById(domId)) getTopAds().then(res => { const data = res.data.data || [] chart.setOption({ title: { text: '热门广告 Top10', left: 20 }, tooltip: { trigger: 'axis' }, xAxis: { type: 'category', data: data.map(item => item.adTitle) }, yAxis: { type: 'value' }, series: [{ type: 'bar', data: data.map(item => item.clickCnt), barWidth: 20 }] }) }) return chart }

图表初始化以后一定要把实例存起来,后续刷新时复用同一个实例。如果每次轮询都调用echarts.init,控制台会警告“There is a chart instance already initialized on the dom”,图表还会越叠越多。

4.4 自动刷新逻辑

大屏的数据不是用户主动刷新出来的,而是要定时轮询。我用setInterval30 秒请求一次,这样页面上的曝光数、点击率会有“跳动感”,演示效果更好:

setInterval(() => { refreshOverview() refreshTrend() refreshTopAds() }, 30000)

需要注意组件销毁时一定要清理定时器,否则路由切换之后,后台还在不断请求接口,白白浪费资源和带宽。这一点做前端联调时最容易踩到。

大屏和后端联调时,最典型的翻车现场是:接口返回的字段是adTitle,前端组件里写的是adName,结果柱状图所有标签都显示 undefined。我建议前后端约定 DTO 字段后,把一份 JSON 示例直接放在接口文档里,前端完全照着字段名取数,顺序乱了也是后端先兜底。

5. 调试记录与避坑清单

5.1 Hadoop 和 Spark 环境搭建时的三个高频问题

我实际搭建时遇到的第一类问题就是 NameNode 启动失败。原因往往是多次执行hdfs namenode -format,导致 VERSION 文件里的 clusterID 不一致。解决办法是把 HDFS 数据目录清空重新格式化,或者直接删除dfs.namenode.name.dir下的目录。但注意,这不是一个可以随便重复执行的操作,一旦集群里有真实数据,重新格式化等于数据全丢,所以生产环境下没人敢乱做。

第二类是 SSH 免密问题。Spark 提交任务到 YARN 时,需要各节点之间能互相 SSH 登录。如果配置不当,作业会一直卡在ACCEPTED状态,看起来像是挂了。检查yarn logs -applicationId <appId>能看到权限相关的报错,解决方法是把公钥分发到各节点。

第三类是spark-submit提示Could not find or load main class。这种情况要么是--class参数写错了全限定类名,要么是 jar 包没有包含依赖。我用maven-shade-plugin生成 fat jar,避免 Spark 作业引用的第三方类缺失。

5.2 Spark 作业运行效率和内存调优

Spark 作业跑挂了,百分之七八十是内存问题。常见报错有这么几类:

报错原因处理方式
OutOfMemoryError: Java heap spaceExecutor 内存不足调大--executor-memory
Container killed by YARN for exceeding memory limits总内存超限调小内存或调大spark.memory.overhead
Shuffle file cannot findExecutor 动态释放关闭动态分配
数据倾斜导致某个 Task 长时间不结束热点 key 集中加盐两阶段聚合或过滤热点广告

最典型的是数据倾斜。广告场景下,一条热门广告可能占了全量点击的 30%,按adId聚合时,这部分数据全部进到同一个分区,对应 Task 就要处理全集群最多的数据,其他 Task 都空闲。我用了一个很土但有效的办法:对adId拼接随机后缀,先做第一轮聚合,再按真实adId做第二轮聚合。这样压力就分散到了多个分区。

内存和线程问题的排查工具,主要依赖 YARN 的 ResourceManager Web UI 和 Spark HistoryServer。这两个界面能清楚看到每个 Executor 的 GC 时间、Shuffle 读写量、Task 处理耗时,比靠猜高效得多。

小节参数建议,我通常这样提交 Spark 作业:

spark-submit \ --master yarn \ --deploy-mode client \ --num-executors 2 \ --executor-memory 2g \ --executor-cores 2 \ --driver-memory 1g \ --conf spark.sql.shuffle.partitions=100 \ --class com.lan.etl.AdEtlJob \ ad-etl-1.0.jar

课程设计阶段数据量不大,两个 Executor 足够。如果后续数据量上来了,优先加 Executor 数量而不是单 Executor 内存,因为单个 Executor 内存在云服务器上很容易把整机压垮。

5.3 Spring Boot 和前端联调问题

Spring Boot 接口写好后,最常见的问题是前端访问报 404 或者跨域错误。404 的原因一般是接口路径写错,尤其是我这种统一的/api/**前缀,前端很容易多写一层或者少写一层。另一个问题是 Redis 缓存反序列化,只要发现接口返回的 JSON 里出现乱码或者ClassCastException,第一反应就应该是查序列化配置。

时间格式问题也很隐蔽。Spark 里用的yyyy-MM-dd HH:mm:ss,到了 Spring Boot 的 LocalDateTime 默认序列化会变成带 T 的 ISO 格式,前端拿到之后直接塞进图表 xAxis 可能没法正常显示。建议用 Jackson 全局配置统一为yyyy-MM-dd HH:mm:ss。

5.4 可视化大屏数据对不上怎么办

大屏数据不一致,通常不是前端或后端的问题,而是上游数据清洗的问题。最典型的例子:点击率指标在 Spark 清洗任务里是按“曝光去重后的点击数/曝光数”算的,但大屏接口单独写了一条 SQL,直接count(click) / count(show),中间漏了去重条件,数字自然就对不上。

碰到这种问题,我都是反着查:先拿接口返回的原始 SQL 在 MySQL 客户端跑一遍,确认 SQL 结果对不对;再拿 Spark 清洗后的 Parquet 文件跑一个聚合,对比两个数。如果两边不一致,优先看过滤条件和去重字段。等两个数一致了,回到前端刷新,问题基本就消失了。

5.5 常见问题速查表

我把这个项目里容易踩的坑统一整理成一张表,开发时遇到问题直接对号入座:

症状原因解决方案
NameNode 起不来clusterID 不一致清空数据目录重新格式化
Spark 作业一直 ACCEPTEDYARN 资源不足或 SSH 不通检查各节点内存和免密配置
推荐列表全空ALScoldStartStrategy未设置设置coldStartStrategy="drop"
Redis 数据乱码序列化配置不对使用 Jackson 或 Fastjson 序列化
大屏柱状图 undefined前后端字段名不一致统一 DTO 字段
前端跨域报错CORS 未配置添加CorsConfig
jar 包启动报 NoClassDefFoundError依赖没有打包使用 shade 插件打 fat jar
数据倾斜跑不动热点 key 过高加盐两阶段聚合

6. 文档编写与源码交付

6.1 技术文档写好,胜过答辩十页 PPT

源码、文档、调试、可视化大屏这四个维度里,文档是最容易被忽略但回报率最高的部分。我的体会是,文档不只是“给别人交付用的”,更是给自己后面重新搭环境用的。这个项目我沉淀了五份文档:

  • 部署文档.md:环境版本、安装步骤、启动命令、常见报错。
  • 接口文档.md:每个接口的 URL、入参、出参、示例 JSON。
  • 数据字典.md:所有表结构和关键字段含义。
  • 测试报告.md:包含推荐效果抽样、大屏指标核对结果。
  • 答辩问题整理.md:针对课程设计或毕设常见问题的准备材料。

README 里则放一份最精简的“启动顺序思维导图”:

启动 HDFS -> 启动 YARN -> 提交 Spark 清洗任务 -> 提交 Spark 推荐模型任务 -> 确认 Redis/MySQL 有数据 -> 启动 Spring Boot -> 启动大屏前端

写文档的时候不要追求文笔,要追求“照着做一定跑得通”。我习惯在装好环境之后,把每一步命令原封不动贴进文档,包括路径和参数。这样读者不需要自己去猜变量,踩坑概率会低非常多。

6.2 源码交付要避开的坑

交付源码时,最尴尬的是把target、node_modules、.idea、logs这些目录一起打包,结果压缩包几十兆,真正有用的代码只有几百 KB。我的整理标准是:所有自动生成的东西都不交付,所有带真实密码的配置都不交付。

项目里我分了三个配置环境:

application.yml # 公共配置 application-dev.yml # 本地开发环境,可保留测试账号 application-prod.yml # 生产环境,密码用占位符替换

数据库初始化 SQL 单独放进data/sql目录,不要跟代码逻辑混在一起。这样别人拿到项目后,第一步建库,第二步改连接串,第三步启动,不会因为少了一张表而卡住。

另外,我用 Git 做版本管理时,.gitignore一定包括这些:

target/ node_modules/ logs/ *.log .idea/ *.iml .DS_Store

6.3 一键启动脚本,减少重复劳动

为了避免每次打开电脑都要敲五六条命令,我在scripts目录下写了一个启动脚本:

#!/bin/bash # 1. 启动 Hadoop start-dfs.sh start-yarn.sh # 2. 提交 Spark 清洗任务 spark-submit \ --master yarn \ --deploy-mode client \ --class com.lan.etl.AdEtlJob \ ad-etl-1.0.jar # 3. 启动 Spring Boot nohup java -Xms512m -Xmx1024m -jar ad-api.jar > api.log 2>&1 & # 4. 启动前端大屏 cd ad-dashboard npm run dev

脚本里每一条命令都是独立的,我是故意没有加set -e。因为 Spark 任务偶尔会失败,但失败后仍然希望 Spring Boot 能启动,API 还能提供上一轮缓存数据。如果你加了set -e,整个脚本会在第一步报错时直接退出,反而弄巧成拙。

最后说点我做这个项目的真实体会

这个项目我前前后后拉了接近两周。最大的感触不是算法难,也不是 Spring Boot 难,而是大数据组件版本之间的兼容性太折腾人。如果重新做一遍,我会先把最小闭环跑通:本地写一个 1000 条的模拟日志文件,用 Spark 读进去,过滤出点击记录,简单聚合出点击次数,然后输出到 Redis,最后在 Spring Boot 里写一个接口手动验证。这个闭环跑通之后,再逐步把 Hadoop、YARN、ALS、可视化大屏一层层加进来。先小步快跑,再扩展枝叶,看起来慢,实际是最快的方式。

如果这个项目后续还要优化,我建议往两个方向发力:一是把 Spark 离线任务改成 Structured Streaming,对用户的新点击行为做分钟级增量更新;二是把推荐结果从简单的 ALS TopN 升级成“规则召回 + 模型排序”两层架构,在排序层加入点击率预估模型。把这些做完,它就不是一个课程设计,而是一个可以撑住真实投放入口流量的在线广告推荐系统了。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/30 15:05:03

从被动挨打到主动防御:专业的邮件安全网关如何重塑企业邮件安全边界

数字化时代&#xff0c;邮件仍是企业对外通信、合同流转、业务协同的核心枢纽。它承载身份、信任与数据资产&#xff0c;也天然成为攻击者最常利用的入口。垃圾邮件、钓鱼攻击、病毒附件、身份伪造、敏感外泄与业务扰动&#xff0c;并非孤立事件&#xff0c;而是攻击链上的不同…

作者头像 李华
网站建设 2026/9/30 15:02:50

SQLite 基本命令与 C/C++ 接口实战:从嵌入式场景到代码实现

SQLite 这个数据库&#xff0c;在嵌入式圈子里基本就是"标配"般的存在。我最早接触它是在做一个车载数据记录仪的项目&#xff0c;内存只有几十兆&#xff0c;却要实时存储 GPS 轨迹、传感器日志、设备状态&#xff0c;还要支持事后按时间范围查询。当时团队里有人提…

作者头像 李华
网站建设 2026/9/30 15:01:50

Smartbits600 测试实战:从开箱到 RFC 2544 吞吐量测试全流程

简介&#xff1a;Smartbits600测试使用指导书是一份面向网络测试初学者与运维人员的实操型文档&#xff0c;围绕NetCom System出品的便携式网络性能测试仪展开&#xff0c;帮助读者从零掌握设备操作与常见测试流程。资源包内共1个doc文件&#xff0c;约977KB&#xff0c;内容按…

作者头像 李华
网站建设 2026/9/30 14:45:18

AST 安全求值

AST 安全求值指的是&#xff1a;把表达式/代码先解析成抽象语法树&#xff08;AST&#xff09;&#xff0c;然后不直接 eval / compile 执行&#xff0c;而是自己遍历 AST&#xff0c;只允许白名单内的节点&#xff0c;并按预定语义解释执行。核心目标是避免任意代码执行、沙箱…

作者头像 李华
网站建设 2026/9/30 14:39:36

【数据集】分省及地级市城投债信用利差数据集(2011-2026年)

数据简介&#xff1a;城投债分省份、地级市信用利差跟踪包括公募债、私募债数据库据库&#xff0c;信用利差个券估值-同期限国开债收益率。剔除剩余期限半年以内或五年以上的个券&#xff0c;估值采用不行权估值&#xff0c;匹配同期限国开债采用插值法。在债券市场中&#xff…

作者头像 李华
网站建设 2026/9/30 14:39:08

提示微调(Prompt Tuning/Prefix Tuning/P-Tuning)技术总结

提示微调属于参数高效微调 PEFT&#xff0c;核心思路&#xff1a;冻结大模型全部主干权重&#xff0c;只训练少量可学习的软提示向量&#xff0c;相比全参数微调显存开销、训练成本大幅下降&#xff0c;是现在大模型落地最常用的方案之一。 1、 核心技术原理提示微调的关键是将…

作者头像 李华