说句实话,第一次看到“基于 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 在线广告推荐的完整链路
一个可复现的在线广告推荐系统,从数据产生到最终展示,完整链路大致是下面这七步:
- 前端投放页面或 SDK 埋点,上报广告曝光、点击、浏览时长等行为日志。
- 日志通过 Nginx 或 Flume 落到 HDFS 的指定目录,按日期分区存放。
- Spark 定时任务读取当日和历史 HDFS 日志,完成数据清洗、字段规整、用户画像聚合。
- 在清洗后的数据上训练或更新推荐模型,生成每个用户的广告 TopN 列表。
- 把推荐结果写入 Redis 或 MySQL,方便后端服务快速读取。
- Spring Boot 对外提供推荐接口和大屏统计接口,前端发起请求获取数据。
- 可视化大屏定时刷新,展示曝光量、点击率、消耗金额、热门广告排行等指标。
很多初学者容易把注意力全放在“算法”上,但实际项目中,链路通畅度比模型复杂度重要得多。一个基于统计规则或协同过滤的简单模型,只要能端到端跑通,就已经超过八成停留在 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 版本不兼容这类问题,项目还没开始写就已经劝退。这个项目我最终锁定的版本组合如下:
| 组件 | 推荐版本 | 说明 |
|---|---|---|
| JDK | 1.8 或 11 | 和 Spark 2/3 生态兼容性最好 |
| Hadoop | 3.3.6 | 内置 NameNode HA 支持,稳定 |
| Spark | 3.3.4 | 支持 PySpark 和 Scala,读 JSON 方便 |
| Spring Boot | 2.7.18 | 避免 3.x 的 Jakarta 迁移问题 |
| MySQL | 8.0 | 存广告元数据和统计结果 |
| Redis | 7.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 格式,每个对象代表一条曝光或点击记录,核心字段如下:
| 字段 | 示例值 | 含义 |
|---|---|---|
| userId | U10001 | 用户唯一标识 |
| adId | A30021 | 广告唯一标识 |
| advertiserId | ADV88 | 广告主标识 |
| categoryId | C_IT | 广告所属类目 |
| channelId | 101 | 投放渠道 |
| deviceType | android | 设备类型 |
| province | 广东 | 地域 |
| isClick | 1 | 是否点击,0/1 |
| showTime | 2025-01-01 10:23:45 | 曝光时间 |
| clickTime | 2025-01-01 10:24:02 | 点击时间 |
| staySeconds | 18 | 落地页停留时长 |
| adTitle | 限时抢购 | 广告标题 |
| token | xxxxxx | 埋点去重标识 |
日志格式一旦定下来,就尽量别改。你可以加字段,但不要调整已有字段的顺序或类型,因为下游 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 space | Executor 内存不足 | 调大--executor-memory |
Container killed by YARN for exceeding memory limits | 总内存超限 | 调小内存或调大spark.memory.overhead |
Shuffle file cannot find | Executor 动态释放 | 关闭动态分配 |
| 数据倾斜导致某个 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 作业一直 ACCEPTED | YARN 资源不足或 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_Store6.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 升级成“规则召回 + 模型排序”两层架构,在排序层加入点击率预估模型。把这些做完,它就不是一个课程设计,而是一个可以撑住真实投放入口流量的在线广告推荐系统了。