做毕设选了“大数据抖音短视频数据分析与可视化”这个方向的同学,我先把话说在前面:这个选题放在今天看,依然是一个性价比很高的选择。它把大数据生态里最常被面试官问到的几个组件(Hadoop、Spark、Hive、Flume、Kafka、ECharts)串成了一条完整的业务链路,同时又自带流量话题属性,做出来的东西视觉效果好,答辩时也容易讲出东西来。本文就把我完成这个项目的完整思路、技术选型、踩坑记录和源码结构全部摊开来讲,适合正在做毕设的本科生、想转行大数据开发的学习者,以及打算拿这个项目当简历亮点的同学参考。
1. 项目定位与整体技术栈设计
1.1 这个毕设到底在做什么:需求拆解
先别急着写代码,把需求理清楚比什么都重要。很多同学拿到这种题目就开始搜爬虫代码,但“抖音短视频数据分析与可视化”这个题目真正的核心,不是爬数据本身,而是“分析”和“可视化”这两个词。
展开来说,这个项目的完整目标应该包括四层:第一,能够拿到一批真实的抖音短视频数据,包含视频的基本信息、作者信息、互动数据、话题标签等;第二,对这批数据做清洗和转换,存到一个适合做离线分析的存储引擎里;第三,用大数据计算框架对数据做多维度统计,比如播放量分布、点赞评论比、发布时间规律、地域热度、话题传播效果等;第四,把统计结果用可视化大屏的方式展示出来,做到“打开页面就能看懂整批数据的情况”。
这四层对应到技术栈上,就是数据采集、数据存储、数据计算、数据展示。你做毕设的时候,一定要在开题报告里把这个链路画清楚,哪怕实现得简单一点,也比只写一个爬虫加几个图表要有说服力得多。
另一个容易犯的错是把项目做得“大而全”,什么都想上。比如有的同学非要加上实时计算、推荐算法、用户画像,结果做完数据采集就已经耗掉大半时间。我个人的建议是,作为毕设项目,先保证离线分析这条线完全跑通、可视化效果好、代码逻辑清晰,有余力再往上加实时模块。因为面试官和答辩老师看重的,是你对每个环节是否真正理解,而不是你机器的服务数量。
1.2 技术选型:为什么是这套组合
技术选型这部分,我直接给出我最终使用的方案,再解释每个组件在这个项目里的具体职责。
| 层次 | 技术组件 | 在项目中的职责 |
|---|---|---|
| 数据采集 | Python + Requests / 抖音开放平台接口 | 获取视频信息、作者信息、互动数据 |
| 消息缓冲 | Kafka | 缓解采集速率与写入速率不匹配的问题 |
| 数据存储 | HDFS | 原始数据与清洗后数据的分布式存储 |
| 元数据管理 | MySQL | 存储最终统计结果,供可视化层查询 |
| 缓存加速 | Redis | 缓存高频查询的统计结果,提升大屏响应速度 |
| 离线计算 | Spark SQL / Hive | 对原始数据进行ETL清洗和多维度聚合统计 |
| 任务调度 | Crontab / 简易调度脚本 | 定时触发采集与统计任务 |
| 可视化展示 | Spring Boot + ECharts + Vue | 后端提供JSON接口,前端渲染大屏 |
这套组合里最容易被质疑的就是“为什么用了Kafka和HDFS,数据量看起来并不大”。我的回答是:这个项目本身是一个教学性质的完整大数据闭环,组件选择的依据不是当前数据量,而是这套架构在面对更大规模数据时的扩展方式。这个回答在答辩时非常加分,因为说明你不是在用技术秀肌肉,而是真正理解了每一层存在的意义。
你需要注意的是Java环境版本和组件版本要提前统一,这是新手最头疼的问题。我这边使用的是JDK 1.8、Hadoop 3.1.3、Spark 2.4.8、Kafka 2.12-2.2.0,这几个版本之间的兼容性测试过,直接照抄能省掉很多麻烦。当然你用更新一点的3.x组合也没问题,但一定不要混用大版本,否则会遇到各种奇怪的序列化和通信报错。
1.3 系统整体架构与数据流向
在写代码前,我强烈建议你先在文档里画出系统的数据流向图。不需要多精美,但要清晰地表达出数据从哪来、经过哪些处理、最终到哪里去。我这里用文字描述一下我的设计,你可以照着这个思路画成自己的图:
采集端Python脚本运行在单独的服务器上,按照配置好的关键词和话题,定时从数据源获取视频列表和详情信息,封装成JSON消息发送到Kafka的douyin-video-topic。然后由一个消费程序从Kafka拉取消息,经过简单的格式校验后写入HDFS的/raw/douyin_video目录,按日期分区。接着Spark任务每天凌晨读取前一天的数据,完成清洗去重、字段解析、维度补全,然后按照需求聚合出各类统计结果,写入MySQL的结果表,同时把部分热点数据写入Redis。最后后端服务查询MySQL和Redis中的数据,组装成JSON交给前端大屏渲染。
这条链路每一步都有明确输入输出,出现问题的时候可以快速定位是哪个环节出的错。我当时第一次跑通全流程的时候,真是长出一口气,因为每层之间的数据传递有很多隐蔽的坑,比如时间格式不统一、字段缺失、JSON解析失败,这些都是要靠日志一点点排查的。
2. 数据采集与预处理:从无到有构建数据集
2.1 数据采集方案:接口选择与合规边界
数据采集是这个项目里第一个绕不开的关卡。很多同学第一反应是写爬虫,但我必须提醒一点:抖音的风控比较严格,单纯用爬虫去抓取页面数据,很容易触发验证码和IP封锁,而且从合规角度讲,未经授权大规模抓取用户数据也存在风险。所以我在做这个项目时,采用的是“开放接口优先、爬虫兜底”的策略。
最稳妥的数据来源是抖音开放平台提供的API接口,通过申请应用获得access_token之后,可以按照官方文档调用指定的接口获取公开数据。虽然开放接口能拿到的字段不如网页端那么丰富,但对于做数据分析来说,核心的播放量、点赞数、评论数、转发数、发布时间、作者信息、视频描述这些字段已经足够了。
如果因为某些原因你无法申请到接口权限,也可以考虑通过合法途径获取公开数据集,比如一些开源的数据社区会有脱敏后的短视频数据包,拿来直接练手完全没问题。我的建议是不要在这上面花太多时间跟反爬机制较劲,一是效率低,二是容易把项目引到灰色地带,没必要。
不管数据来源是哪里,你都需要在采集端做好三件事:一是字段统一,无论来源格式是JSON还是别的结构,落到Kafka之前都转成统一结构,方便后续处理;二是异常处理,网络超时、返回码异常都要有重试和跳过机制,保证采集任务能长期稳定运行;三是采集频率控制,设置合理的间隔时间,既不影响目标站点,也避免把自己机器搞崩。
2.2 数据清洗与ETL的实战细节
数据采回来只是第一步,真正的体力活在于清洗。我第一批采集的数据里有将近三成是有问题的,包括但不限于:字段为空、播放量显示为“w+”(比如1.2w这种格式)、JSON嵌套层级不一致、时间戳和字符串格式混用、标签字段是数组还是字符串不统一等。这些脏数据如果不处理,直接拿去聚合,出来的统计结果根本没法看。
我使用的是Spark SQL做ETL,核心步骤大致是:读取HDFS原始目录的JSON文件,用from_json函数配合schema定义把嵌套JSON解析成结构化列;然后做去重,以视频ID为唯一键,保留最新一条记录;接着做字段类型转换,把“1.2w”这种格式统一转成数值型,这里需要自己写一个UDF,处理中文单位换算;再做时间标准化,统一成yyyy-MM-dd HH:mm:ss格式;最后过滤掉核心字段为空的数据,写入清洗后的HDFS目录。
ETL过程中有一个细节很多人忽略:数据质量问题要在源头控制,而不是在分析阶段补救。比如播放量字段如果允许为空,后续做排名、分组、求均值的时候都会出现逻辑错误。我建议你在ETL阶段就定义好每张表的字段完整性规则,宁可丢弃少量脏数据,也不要让它们污染后续所有统计结果。
清洗完的数据建议注册成Spark临时表,方便后续分析和校验。你可以先跑几个简单的select count(*)和select * limit 10,确认数据形态符合预期再往下一步走。这一步虽然枯燥,但能帮你省掉后面大量排查问题的时间。
2.3 数据存储设计:HDFS + MySQL + Redis的分工
这个项目的存储设计有一个明显分层:HDFS存原始和清洗后的数据,MySQL存聚合后的指标结果,Redis存热点查询的缓存。
很多同学会问,这么点数据量有必要用HDFS吗?我的答案是,从纯功能角度确实不是必须的,但从毕设的技术完整性和学习价值角度,HDFS是必不可少的。因为你要学习的是分布式文件系统的思维方式和文件组织方式。实际操作中,你会在HDFS上建立/raw和/clean两个根目录,下面按业务和时间分区,比如/raw/douyin_video/date=2024-01-01/。这种按照日期分区的组织方式,能够让后续的Spark任务很方便地通过分区裁剪只读取需要的日期范围,避免全表扫描。
MySQL在设计表结构时,我建议把结果表按照分析维度拆开,而不是把上百个指标塞进一张大宽表。比如视频表现分析表、作者榜单表、话题热度表、地域分布表、发布时间分布表各一张,每张表都包含统计维度和指标字段。这样设计的好处是:方便单独更新某张表、接口查询逻辑简单清晰、后续增加新指标不需要改动已有表结构。
Redis在这条链路里的作用是查询加速。大屏页面打开的时候会一次性请求十几个接口,如果每次都去MySQL里查,虽然数据量不大不会慢,但为了演示“缓存层”这个技术点,可以对榜单类结果设置5分钟过期时间缓存到Redis。用户在操作大屏的时候可以明显感觉到切页流畅、响应快,答辩时也能顺带讲出“缓存的使用解决了什么问题”这个技术话题。
3. 核心分析指标与Spark离线统计实现
3.1 分析维度设计:从5个层面拆解指标
可视化大屏好看的前提是分析维度有逻辑,不是随便拼几个图表上去。我做这个项目时把指标拆成五个维度的模块,每个模块对应大屏上的一块区域。
第一是内容生态概览,包括视频总量、参与作者数、总播放量、总点赞量、总评论量、总转发量这些总体指标。第二是视频表现分析,包括播放量分布区间、点赞评论转化率、完播率估算、互动率排名等。第三是作者维度分析,包括作者活跃度排行、高产作者Top10、互动效果最好的作者Top10。第四是内容特征分析,包括话题标签的使用频率、不同视频时长的表现差异、视频描述的关键词分布。第五是时间与地域分析,包括视频发布时间的小时分布、星期分布、作者所在省份/城市的热度排名。
这个维度拆解不是随便拍的,每一层都有明确的分析目的,比如“时间发布规律”可以侧面反映短视频用户的使用高峰时段,“地域分布”可以反映不同地区的创作者活跃情况,“互动率排名”能够筛选出真正有内容影响力的账号,而不是只看绝对播放量。
有了完整指标矩阵,你在设计Spark统计任务时心里才有一个全局图。把指标和SQL的对应关系列出来,比如“互动率=点赞数+评论数+转发数 除以 播放量”,先定义清楚再动手写代码,会少走很多弯路。
3.2 Spark SQL统计任务的配置与调优
Spark SQL是这套离线统计的核心计算引擎。最基本的用法是读取清洗后的数据,注册成临时表,然后写各种SQL做聚合。
我以“统计每个省份的平均播放量Top10”为例来说明整个写法和思路。大致步骤如下:用Spark读取HDFS清洗目录,得到DataFrame后注册为video_table;然后创建省份与视频的关联表province_table;再执行类似下面的SQL:
SELECT province, ROUND(AVG(play_count), 2) AS avg_play, COUNT(*) AS video_cnt FROM video_table WHERE province IS NOT NULL AND play_count > 0 GROUP BY province ORDER BY avg_play DESC LIMIT 10这样一句简单的SQL背后其实做了几件重要事情:数据过滤、分组聚合、排序截取。真正的项目中要把所有指标都写成这样的SQL,封装在一个AnalysisRunner对象里,每个方法负责一个模块的统计,输出结果写入不同的MySQL表。
关于Spark任务性能调优,有两个参数值得提前设置好。第一个是spark.sql.shuffle.partitions,默认是200,对于这个项目的数据量来说太大了,建议改成10到20,否则每个reduce任务都要落盘,浪费调度开销。第二个是spark.sql.adaptive.enabled,如果使用Spark 3.x版本,可以打开自适应查询执行,让Spark根据数据量自动调整分区数。
还有一个经验是尽量复用同一个SparkSession。如果你对每个指标都重新创建一个SparkSession,启动开销就能让你多等十几分钟。把Session初始化放到程序入口,所有统计方法都使用同一个Session实例,整个任务跑完大约比逐个创建快一倍以上。
3.3 结果表设计与更新策略
统计结果要落到MySQL,表结构的设计得提前想好。我给出几个核心表的字段设计,你按需要扩展即可。
视频表现分析表:
| 字段名 | 类型 | 说明 |
|---|---|---|
| id | int | 自增主键 |
| video_id | varchar | 视频ID |
| title | varchar | 视频标题 |
| play_count | bigint | 播放量 |
| like_count | bigint | 点赞量 |
| comment_count | bigint | 评论量 |
| share_count | bigint | 转发量 |
| interact_rate | decimal | 互动率 |
| duration | int | 视频时长秒数 |
| publish_date | date | 发布日期 |
作者榜单表类似,核心字段是author_id、author_name、follower_count、total_play、total_like、video_count、avg_interact_rate等。地域分布表的核心字段是province、city、video_count、total_play、total_like等。
更新策略上我采用的是“先删后插”的简单方式:每个统计任务跑完后,用DELETE FROM table WHERE stat_date = '前一天日期'清掉对应分区的旧数据,然后批量插入新结果。数据量本身不大,这种粗暴的方式反而简单、可靠,也可以很自然地保证幂等性。
批量插入的时候我强烈建议你使用JDBC的addBatch和executeBatch,而不是一条一条insert。不要小看这个细节,同样写一万条结果,批量插入只需要几秒,逐条插入可能要几分钟,而且对MySQL的压力也大得多。
4. 可视化大屏:从数据到图表的最后一公里
4.1 大屏布局与图表选型
可视化大屏是整个项目最直观的成果,不管数据分析做得再好,大屏效果不行,答辩印象分会直接打折扣。我在设计大屏布局时参考了主流数据产品的风格:深色背景、亮色主题、左中右三栏布局。
具体的区块规划是:顶部是标题栏和核心KPI数字,比如总视频数、总播放量、总点赞数,这些数字用大号字体配合翻牌动画显示。左侧从上到下放“作者活跃度Top10”排行榜柱状图和“话题热度Top10”横向条形图。中间部分放“播放量趋势折线图”和“视频时长分布饼图”。右侧放“地域分布地图”和“发布时间24小时分布柱状图”。底部可以放一个滚动的数据列表,展示最新采集到的视频信息。
图表选型上我使用的是ECharts,生态成熟、文档丰富、图表类型齐全。地图部分需要使用中国地图的GeoJSON注册,ECharts 5.x版本里已经不再内置地图数据,需要额外引入或下载中国地图的JSON文件。省份名称要与数据表中的省份维度完全一致,否则地图渲染会出现空白,这是被很多人忽略的一个细节问题。
4.2 前后端接口设计与实时刷新
大屏页面和后端接口之间采用JSON通信。后端用Spring Boot构建RESTful接口,每个模块对应一个接口路径,比如/api/overview返回总览KPI,/api/author/top10返回作者榜单,/api/hotwords返回热点标签等。前端用Axios请求这些接口,拿到数据后填充到ECharts的配置项中。
关于“实时刷新”,这里需要区分一下:抖音短视频数据本身是离线采集和分析的,所以大屏并不需要做到秒级真实流式刷新。我的做法是前端每60秒向后端拉取一次最新数据,后端先查Redis缓存,有就直接返回,没有就去MySQL查询并回填Redis,这样在演示时给人“数据在自动更新”的感觉。
如果你想让项目显得更有技术含量,可以在采集端模拟“实时上报”的逻辑:每次采集到一条新数据,先发到Kafka,然后有一个消费者把最新的一条数据写入Redis的列表里。大屏轮询的时候同时读取Redis中的最新数据和MySQL中的聚合结果,两个数据源拼接后返回给前端展示。这样也能让面试官感受到你确实理解流式处理的场景。
4.3 大屏部署的几个关键细节
大屏项目在本地开发环境跑得好好的,一部署到服务器上各种问题就出来了。这里分享几个部署阶段最常踩的坑。
第一是跨域问题。前端页面如果和后端接口不在同一个域名下,浏览器会拦截跨域请求。解决办法可以是在后端加@CrossOrigin注解,也可以使用Nginx反向代理统一入口。推荐用Nginx,因为最终部署时前端构建产物和后端服务往往不在一台机器上,Nginx可以同时代理静态页面和后端API路径,一举两得。
第二是ECharts地图资源的路径问题。地图JSON文件放在前端项目里,构建时如果路径不对会导致请求404。建议把地图JSON放到静态资源目录下,并使用相对路径引入,避免部署到子目录时找不到资源。
第三是大屏在不同分辨率下的适配问题。答辩现场往往用的是投影仪或者大显示器,建议使用屏幕宽度自适应方案。具体做法是把大屏整体设计在1920x1080的基准分辨率里,再用CSS的transform: scale()对根容器做等比缩放,这样无论现场屏幕是不是16:9,大屏都不会因为拉伸导致布局错位。
5. 源码结构梳理与复现指南
5.1 源码目录结构与核心模块说明
拿到源码或者准备自己搭框架的同学,第一步一定先搞明白目录结构,不然看代码时会一头雾水。我给出一个经过整理的标准目录结构,你可以对照自己项目的源码来理解。
douyin-data-analysis/ ├── conf/ # 配置文件目录 │ ├── application.yml # Spring Boot配置 │ ├── spark-conf.properties # Spark任务的参数配置 │ └── log4j2.xml # 日志配置 ├── collector/ # 数据采集模块 │ ├── main.py # 采集入口脚本 │ ├── api_client.py # 数据源接口封装 │ └── producer.py # Kafka生产者逻辑 ├── etl/ # 数据清洗模块 │ ├── spark_etl.py # Spark ETL脚本 │ ├── udfs.py # 自定义UDF函数 │ └── schema.py # 数据结构定义 ├── analysis/ # 统计模块 │ ├── SparkAnalysisRunner.scala │ ├── VideoAnalyzer.scala │ ├── AuthorAnalyzer.scala │ └── TopicAnalyzer.scala ├── dashboard/ # 可视化大屏前端 │ ├── src/ │ │ ├── views/Dashboard.vue │ │ └── components/ │ │ ├── KpiCard.vue │ │ ├── AuthorRankChart.vue │ │ └── MapChart.vue │ └── package.json ├── backend/ # 后端服务模块 │ ├── pom.xml │ └── src/main/java/com/example/douyin/ │ ├── controller/ │ ├── service/ │ └── mapper/ └── docs/ # 文档目录 ├── 系统设计文档.md ├── 接口文档.md └── 部署文档.md采集模块负责人任务是跟外部数据源打交道,代码结构上要把接口封装和业务逻辑分开,方便替换数据来源。ETL模块核心是自定义UDF函数和数据清洗流程,这里的代码一定要写注释,因为后续改字段映射时最容易在这里翻车。分析模块是Spark任务的Scala代码,每个分析器对应一个分析维度,输出写MySQL。后端模块是大屏的数据出口,只做查询不做复杂逻辑,代码越简单越不容易出错。前端模块是门面,图表的交互和动画都在这里控制。
5.2 从零复现的七步操作流程
如果你是从零开始搭整套环境,我建议严格按下面这个顺序来,先保证环境能用,再写业务逻辑。
第一步,准备三台虚拟机或者一台内存16G以上的服务器,安装CentOS 7.9,配置好静态IP和主机名。搭建Hadoop完全分布式集群,包括NameNode、DataNode、ResourceManager、NodeManager,这一步建议参考Hadoop官方文档操作。第二步,安装Zookeeper和Kafka集群,Kafka依赖Zookeeper管理元数据,一定要先启动Zookeeper再启动Kafka。第三步,安装Spark,确认能通过spark-shell正常打开交互式环境。第四步,安装MySQL和Redis,创建项目需要的数据库和用户,初始化表结构。第五步,把采集模块部署到一台机器上,先运行一次单条采集验证整条链路通不通,再开启循环采集。第六步,启动Spark任务做ETL和统计分析,观察日志和MySQL结果表是否有数据。第七步,部署后端服务和前端大屏,浏览器访问页面确认所有图表正常加载。
这些步骤中,第一步和第二步是最耗时间的,往往会出现各种进程启动失败、端口被占用的问题。我建议你按照错误日志倒序排查,也就是先看最后几行报错,再去定位具体配置。大部分情况下都是主机名和IP映射没配好、防火墙没关、配置文件格式不对这三类问题。
5.3 项目答辩时的亮点提炼与扩展方向
这个项目在答辩时,如果只说自己会跑通整个流程,最多算及格。真正拿高分的,是你能够回答清楚每一个技术组件“为什么存在”以及“替换掉行不行”。
我给你列几个必答的亮点问题:第一,为什么用Kafka而不是直接把数据写入HDFS?答案是Kafka起到了削峰填谷的作用,即使采集端爆发式写入,下游写入HDFS的速率是稳定的,不会压垮文件系统。第二,Spark和Hive的定位有什么区别?答案是Hive适合跑定时SQL批处理,Spark提供了更灵活的RDD、DataFrame和SQL统一编程模型,更容易做精细化的数据清洗和复杂指标计算。第三,如果数据量上升到每天上亿条,哪里会成为瓶颈?答案是Kafka的分区数和Spark的并行度,横向扩展的方向是增加Kafka分区数并增加Spark Executor数量。
扩展方向上,你可以考虑在现有基础上加入Flume做日志采集、加DolphinScheduler做可视化任务调度、用ClickHouse替换MySQL做大规模结果查询、增加推荐算法模块分析用户兴趣等。答辩时如果能说出“当前项目已经跑通主链路,后续可以从实时性和智能化两个方向继续演进”这类话,会给老师留下一个思路开阔的印象。
6. 实测踩坑记录与常见问题排查
6.1 环境搭建阶段的经典报错
整套环境搭建过程中,我遇到最多的三类报错,列出来给大家避坑。
第一类是Hadoop启动后DataNode没有正常启动。原因通常是NameNode和DataNode的clusterID不一致,解决办法是删除DataNode目录下的VERSION文件并重新执行hdfs namenode -format。不过这里有个更稳妥的操作顺序:先格式化NameNode,再启动HDFS,最后再启动YARN,不要在启动过程中反复格式化,容易把元数据搞乱。
第二类是Spark连接HDFS时出现的org.apache.hadoop.security.AccessControlException权限问题。排查方向通常是HDFS目录的写入权限,最简单的方式是给运行用户配置对应的目录权限,或者临时设置HDFS的权限为宽松模式。测试环境可以这么做,生产环境还是建议用Kerberos做安全认证,但毕设阶段不必搞那么复杂。
第三类是Kafka生产者能连上Broker,但消费者一直拉不到消息。排查步骤是:先用kafka-topics.sh --describe确认topic存在且分区正常,再用kafka-console-consumer.sh测试能否消费到消息,如果消费端也收不到,重点检查采集脚本是否真正发送成功,看Kafka服务端日志里有没有错误。很多时候是生产者发送时没有设置acks=1导致消息没有真正写入,这一点非常隐蔽。
6.2 数据质量和分析结果异常的处理
数据质量问题是离线分析过程中最常见的坑。我遇到过播放量字段出现null、重复记录、同一视频在不同时间被采集了多次且数据不一致等问题。
处理重复记录,我是在ETL阶段用row_number() over (partition by video_id order by collected_at desc)取最新一条,保证每个视频ID唯一。处理null值,要先明确这个字段后续是否参与计算,如果参与计算就丢弃或填充默认值,如果不参与就保留。处理数据不一致,比如同一视频上一次采集播放量是1000,下一次变成800,要考虑是否采集到不同状态的数据,这种情况我选择以最后一次采集为准。
还有一类问题是统计结果对不上常识。比如地域分布里某个省份的视频数量异常高,排查后发现是采集脚本中省份字段解析有误。我的经验是每个分析结果都要能往回溯源,从最终图表一层层查到原始数据,这样才能快速定位是统计逻辑问题还是原始数据问题。
如果你发现自己清洗后的数据量和原始数据量差别很大,不要急着怀疑代码,先检查源数据本身的质量。我之前采集的一批数据里近四成是空壳账号发布的视频,没有播放量也没有点赞数,这种数据即使清洗后对分析也没有价值,直接过滤掉反而是合理的。
6.3 性能优化与后续演进方向
项目做完之后的性能优化,可以从三个方向入手。
第一个方向是提高Spark统计任务的执行效率。核心思路是减少Shuffle数据量。比如关联操作时,如果一张表很小,用Broadcast Join广播小表,避免Shuffle。再比如尽量在分组前先过滤掉不需要的数据,而不是全量聚合后再过滤。这些优化对当前数据量提升不明显,但能体现你对Spark原理的理解。
第二个方向是增加采集端的稳定性和可观测性。可以给采集任务加上心跳上报和失败告警,配合Azkaban或DolphinScheduler做任务失败自动重跑。这样整个链路从采集到统计都不需要人工干预,更像一个真正在运行的数据平台。
第三个方向是我个人觉得最有意思的升级思路:在现有的离线统计基础上,增加一个“热门视频实时排行”模块。用Spark Streaming或者Flink消费Kafka中的新数据,每5分钟更新一次热门视频榜单,写入Redis。这样大屏上除了离线指标,还能展示“当前正在上升的热门内容”,整个项目的实时感和科技感瞬间拉满。
关于这个项目的源码复现,一个实际心得是:不要试图一次性把整套环境部署完再开始写业务代码,而是采用“最小闭环验证法”,先用一份很小的测试数据把采集、存储、统计、展示整条链路跑通,确认每个环节的输入输出都符合预期,再扩充数据量和分析维度。这样调试周期短,出问题也容易定位。
如果你打算把这个项目作为简历上的核心项目,建议重点讲清楚两个故事:一是你在数据质量和ETL层面做了哪些细致工作,二是你是如何设计指标体系来支撑可视化决策的。这两块内容面试官普遍比较感兴趣,因为很多人做类似项目都停留在“能跑通”的层面,而你能说出“为什么这么做”,这就是差异点。