简介:这是一份面向本科毕业设计及大数据入门学习者的完整项目源码包,围绕新闻网站用户浏览日志,提供从Flume日志采集、HBase存储、Spark Streaming实时消费到可视化展示的全链路实现。项目基于Spark2.x,整合Hive、Kafka、Grafana等组件,覆盖用户行为采集、实时流处理、离线分析与指标展示等核心环节,可用于毕设系统搭建、实时计算学习或大数据课程设计参考。资源共35个文件,压缩包约3.46MB,类型以Scala/Java源码、JAR依赖、XML/Maven配置为主,另有MD项目说明、TXT参考步骤、PNG可视化截图及HTML/JS前端展示文件。目录按功能划分,包括Flume与HBase集成示例、Spark测试脚本、Spark Streaming业务处理代码,便于对照阅读和二次开发。项目可实现实时统计访问量前20的新闻话题、曝光话题及分时段浏览量峰值等指标,体现实时与离线结合的场景化设计。已有68人学习下载,配套文档和参考步骤可帮助快速复现环境,适合需要参考完整大数据分析项目结构的读者。
1. 基于 Spark2 的新闻浏览日志实时分析与可视化:这个毕设题值不值得做
要做好“基于Spark2的新闻浏览日志大数据实时分析与可视化系统”,本质不是写一个 Spark WordCount,而是把“日志产生 → 消息队列 → 流式计算 → 结果存储 → 大屏展示”这条实时链路完整跑通。我见过不少同学把这题做成离线统计,最后答辩被问“实时在哪”就卡住;反过来,如果你能当场演示一次点击在几秒内出现在大屏上,这个毕设就立住了。
它适合两类人:一是想往大数据开发方向走的,这套架构里 Kafka、Spark Streaming、Redis、WebSocket 全是面试高频词;二是时间不算充裕、需要“可运行比炫技更优先”的,这套方案单机 4GB 内存就能跑。下面按架构、数据模拟、实时计算、避坑、可视化的顺序,把可以直接开工的路径和代码给你。
2. 先把数据链路立住:Spark2、Kafka、Redis、ECharts 在这一题里各司其职
拿到压缩包先别急着解压跑代码,把里面那份操作步骤文档先翻一遍,按环境清单把 Kafka、Spark、Redis、Zookeeper 装好。我第一次就是跳过文档直接跑 Spark,结果 Zookeeper 没起,Kafka 一直连接超时,浪费了半天。下面这套架构就是那条链路的地图版,你对着一章就能把组件关系和指标口径都定下来。
2.1 为什么选 Spark2 而不是 Flink 或 Spark3
毕设选型的第一原则是“在答辩能讲清楚的前提下选生态最成熟的”。Spark2 在这一点上有三个优势:第一,教程和博客存量极大,StreamingContext、DStream、reduceByKeyAndWindow 这些关键词随便一搜就是完整例子,遇到问题抄作业都容易;第二,单机伪分布式跑起来很轻,虚拟机分 4GB 内存就能把 Streaming、Kafka、Redis 一起塞进去,Flink 在同样配置下光状态后端和内存调优就要多花几天;第三,Spark2 和 Hadoop 2.x 的兼容性在毕设环境最常见,很多实验室集群就是这套组合。
要澄清一个容易误解的点:Spark Streaming 是微批模型,不是逐条处理。它会先把到达的日志按 5 秒或 10 秒切成一个批次,再交给 Spark 引擎算。对有几十万用户、QPS 在几千到几万之间的新闻站来说,5 秒延迟完全够用。真正需要毫秒级响应的场景才考虑 Flink,而“新闻浏览日志实时分析”这个业务,微批反而更好讲窗口语义。如果指导老师指定 Spark3 或要求用 Structured Streaming,也没关系,后面 4.4 节会给出两种写法的差异,DStream 的作业改成 DataFrame 语法成本很低。
2.2 一条完整的数据流:从日志产生到前端大屏
整套系统可以拆成五段,每一段解决一个明确的问题。第一段是数据接入,用 Kafka 做缓冲削峰,避免日志突发流量直接把计算层打挂;第二段是实时计算,Spark Streaming 消费 Kafka 里的日志,算出 PV、UV、热点新闻 TopN、来源渠道占比这些指标;第三段是结果存储,Redis 承接高频写入,用 Sorted Set 天然支持排行,用 Set 做 UV 去重;第四段是推送,后端通过 WebSocket 把新算出的指标推给浏览器;第五段是展示,ECharts 渲染大屏,setOption 增量更新就能实现图表平滑刷新。
组件选型可以对照下面这张表,答辩时这就是你的架构设计依据:
| 层级 | 组件 | 选型理由 |
|---|---|---|
| 消息队列 | Kafka | 日志场景事实标准,分区模型与 offset 机制天然适配流式计算 |
| 实时计算 | Spark Streaming | 微批模型,5 秒窗口即可满足新闻日志实时性要求 |
| 结果存储 | Redis | 内存读写快,Zset 支持 TopN,Set 支持 UV 去重 |
| 推送 | WebSocket | TCP 长连接,服务端主动推送,避免高频轮询 |
| 展示 | ECharts | 开箱即用,折线图、柱状图、饼图都能直接对接 JSON 数据 |
这里想特别提醒一句:不要让 Spark Streaming 直接操作 MySQL。实时场景下每秒都有窗口结果要写,MySQL 的连接和行锁会成为链路里最先堵死的点。Redis 的正确用法是当“指标缓存层”,Spark 算完写 Redis,后端接口也从 Redis 读,MySQL 只在最后做日终落库或离线报表时才参与。很多毕设源码把自己写成大量 SQL 拼接,本质上还是离线思路。
2.3 指标口径:先想清楚“实时”是什么粒度
动手写代码前,最值得花时间的是把指标口径定死。常见做法是先定义四个核心指标。实时 PV 是每分钟窗口内所有点击行为的总数,用窗口聚合直接累加;实时 UV 是同一分钟内去重后的活跃用户数,去重逻辑基于 user_id,精确到分钟级即可;热点新闻 TopN 是新闻维度 PV 的倒序排行,取前 10 或前 20;来源渠道占比是把 channel 字段按推荐页、分类页、搜索、推送分类计数,用于画饼图。
这里有个毕设答辩必问的点——“实时”到底是多久。建议在文档里明确写:系统采用 5 秒微批粒度,指标按 60 秒滑动窗口统计,窗口每次滑动 5 秒。这意味着大屏上看到的数值最多滞后 5 秒。把这个口径写清楚,答辩时就不容易被人抓住“不是真实时”的话柄。窗口时长和滑动间隔不要拍脑袋定,后面 4.3 节会让你看到这两个参数直接决定 Redis 写入频率和前端刷新频率,三者必须联动。
3. 让数据先流起来:新闻日志模拟器与 Kafka 接入代码
没有真实业务日志的时候,不要干等数据,自己写一个模拟器把 Kafka 灌满是最快的启动方式。这一章解决“数据从哪来”的问题,直接用 Python 脚本生成结构化日志并写入 Kafka,同时把生产端几个关键参数的取舍讲清楚。
3.1 日志长什么样:五个字段决定后续所有指标
这套系统不依赖真实前端埋点,毕设阶段用模拟日志完全足够。日志格式我一般控制在五到六个字段,字段多了模拟器要编的值就多,下游解析代码也会变长,反而增加排查成本。下面的表就是一份可直接用的日志结构,news_id 决定热点排行,user_id 决定 UV,channel 决定渠道占比,ts 决定窗口归属,city 是留着做地图可视化或城市维度扩展的。
| 字段 | 类型 | 含义 | 示例值 |
|---|---|---|---|
| news_id | string | 新闻唯一 ID | news_0123 |
| user_id | string | 用户唯一 ID | user_8899 |
| channel | string | 用户入口渠道 | 推荐页 / 分类页 / 搜索 / 推送 |
| city | string | 城市 | 北京 / 上海 / 广州 |
| ts | long | 事件发生毫秒时间戳 | 1712113456789 |
注意 ts 用毫秒而不是秒,后面 Kafka 和 Spark 里会有好几处时间参与窗口计算,前端展示时除以 1000 转成可读时间即可,跨单位换算的坑能避就避。channel 只保留四个固定值也是刻意为之,为了让渠道占比饼图看起来稳定,随机值太多会导致图几乎每条都变。
3.2 用 Python 模拟器往 Kafka 里灌数据:最小可运行代码
常见做法是写一个 Python 脚本,从预置的 200 条新闻、1 万个用户里随机抽,每隔 10 毫秒发一条点击日志。用 kafka-python 库,连接本机 Kafka,序列化用 JSON。下面这段代码可以直接存成 news_log_simulator.py 运行:
import json import random import time from kafka import KafkaProducer news_pool = [f"news_{i:04d}" for i in range(1, 201)] user_pool = [f"user_{i:04d}" for i in range(1, 10001)] channel_pool = ["推荐页", "分类页", "搜索", "推送"] city_pool = ["北京", "上海", "广州", "深圳", "杭州"] producer = KafkaProducer( bootstrap_servers="localhost:9092", acks=1, retries=3, linger_ms=20, batch_size=65536, value_serializer=lambda v: json.dumps(v).encode("utf-8") ) print("开始模拟新闻浏览日志,每轮 10ms...") while True: log = { "news_id": random.choice(news_pool), "user_id": random.choice(user_pool), "channel": random.choice(channel_pool), "city": random.choice(city_pool), "ts": int(time.time() * 1000) } producer.send("news_log", value=log) producer.flush() time.sleep(0.01)这段代码的逻辑很简单:池子里随机抽值组装一条 JSON,通过 producer.send 发到 news_log 这个 topic,flush 保证当前消息真正发出后再进入下一轮。第一次调试时建议先启动 Kafka 再运行脚本,然后用消费命令确认数据真的进来了:
kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic news_log --from-beginning能看到一行行 JSON 刷屏,说明 Kafka 侧已经通了。这里不用急着做压力测试,模拟器本身要调的参数就是 0.01 秒的 sleep,这就是控制 QPS 的旋钮:sleep(0.01) 大约是每秒 100 条,想模拟高峰改成 sleep(0.001),想低峰改成 sleep(0.1)。毕设答辩时让大屏动起来,每秒 100 条完全够用。
3.3 生产端参数:acks、linger.ms、batch.size 怎么调
上面代码里那四个 KafkaProducer 参数不是随便写的。acks=1 表示消息写入 leader 分区就算成功,兼顾吞吐和可靠性;单机环境追求更强一致可以改成 all,但会拖慢发送速率,毕设场景不值得。linger_ms=20 是“攒批”时间,意思是消息先积压在发送缓冲区最多 20 毫秒再一起发出去,这是吞吐和延迟之间的折中;设成 0 就是来一条发一条,延迟最低但请求数量剧增。batch_size=65536 是发送缓冲区每批的上限字节数,配合 linger_ms 使用,64KB 是比较中庸的值。
Kafka 侧的 topic 也要先建好,分区数建议直接给 3。这样后面 Spark 端每个分区一个 consumer 线程,并行度刚好铺开,不会出现一个分区超载而其他分区空闲的倾斜情况。建 topic 的命令是:
kafka-topics.sh --create --topic news_log --partitions 3 --replication-factor 1 \ --bootstrap-server localhost:9092副本因子设 1 是因为毕设环境通常只有一个 broker,设 2 以上会报“副本数大于 broker 数”的错误。还有一个容易漏的细节:Kafka 默认日志保留 7 天,对毕设演示周期够用,但如果模拟器连续跑几天不关,磁盘会被 topic 日志撑爆,建议在 server.properties 里把 log.retention.hours 调到 24 小时以内,防止演示那天磁盘满了导致消费异常。
4. Spark Streaming 实时统计:PV、UV、热点新闻排行的核心代码与调参经验
数据进了 Kafka,剩下的就是 Spark 的事了。这一章从创建 StreamingContext 开始,到消费 Kafka、窗口聚合、写 Redis,把一条完整可跑的实时计算链路的代码和参数都过一遍。
4.1 创建 StreamingContext:batchDuration 与并行度的第一次取舍
Spark Streaming 的入口是 StreamingContext,两个最关键的参数是 batchDuration 和并行度。batchDuration 就是前面说的 5 秒微批周期,我默认给 5 秒;如果机器只有 2 核内存 4GB,可以放宽到 10 秒,但窗口计算里的滑动间隔也要同步改。并行度在本地直接写 local[4],代表用 4 个线程跑,其中 1 个线程负责接收数据,另外 3 个做窗口计算,避免接收和计算挤在同一个线程里互相拖累。
import org.apache.spark.SparkConf import org.apache.spark.streaming.{Seconds, StreamingContext} val conf = new SparkConf() .setAppName("NewsLogRealtimeAnalytics") .setMaster("local[4]") .set("spark.streaming.backpressure.enabled", "true") .set("spark.streaming.kafka.maxRatePerPartition", "2000") val ssc = new StreamingContext(conf, Seconds(5)) ssc.checkpoint("checkpoint/news-log")这段代码里有三个点要说明。第一,setMaster 只在本地调试时保留,集群提交时删掉这行,用 spark-submit --master yarn 替代,否则集群环境会报 master 冲突。第二,backpressure 和 maxRatePerPartition 是配套用的,backpressure 让 Spark 根据处理速度反过来限制 Kafka 消费速率,maxRatePerPartition 设置每个分区每秒最多消费 2000 条,这是防内存被打爆的保险丝。第三,checkpoint 目录必须设置,它不仅保存 RDD 数据,还保存 offset 和窗口元数据,是第五章里“重启不丢数据”的前提。
4.2 从 Kafka 消费:用 createDirectStream 而不是 createStream
Spark 2 里消费 Kafka 有两套 API,0-8 时代的 KafkaUtils.createStream 用的还是 Receiver 方式,数据先落 WAL 再处理,offset 由 Kafka 自动维护;0-10 引入的 createDirectStream 直接把每个 Kafka 分区映射成一个 RDD 分区,offset 由 Spark 手动管理。毕设和实际生产我都建议用 createDirectStream,因为消费和计算在同一进程里,不会出现“先接收后处理”导致的两段式延迟,offset 也能精确控制。
import org.apache.spark.streaming.kafka010._ import org.apache.kafka.common.serialization.StringDeserializer val kafkaParams = Map( "bootstrap.servers" -> "localhost:9092", "key.deserializer" -> classOf[StringDeserializer], "value.deserializer" -> classOf[StringDeserializer], "group.id" -> "news-realtime-group", "auto.offset.reset" -> "latest", "enable.auto.commit" -> "false" ) val topics = Array("news_log") val stream = KafkaUtils.createDirectStream[String, String]( ssc, LocationStrategies.PreferConsistent, ConsumerStrategies.Subscribe[String, String](topics, kafkaParams) ) val logs = stream.map(record => { val json = new JSONObject(record.value()) (json.getString("news_id"), json.getString("user_id"), json.getString("channel"), json.getLong("ts")) })参数逐一说:bootstrap.servers 填 Kafka 的地址,本地就是 localhost:9092;group.id 是消费者组名,重启后 Spark 能通过它找到上次消费的位置;auto.offset.reset=latest 表示从最新位置开始读,适合演示场景,想重放历史数据改成 earliest;enable.auto.commit=false 是关键,它关闭 Kafka 自动提交 offset,让 Spark 在处理完一批数据后再手动提交,避免“数据没算完就把 offset 提交了”导致丢数据。手动提交的写法放在 5.1 节,那是一个必须单独强调的坑。
消费到 record 后,只取四个字段转成 Tuple4,这一步相当于把 JSON 字符串拆成结构化表示,后续所有窗口计算都在这条 logs 流上做。这里可以先不解析成 case class,因为 DStream 的泛型在 Scala 里用 Tuple 更省事,传参和 map 操作都比较直白。
4.3 核心指标计算:窗口聚合、UV 去重、TopN 落地 Redis
PV 和热点排行用同一个窗口算子 reduceByKeyAndWindow。它比简单的 reduceByKey 高级在能维护一个滑动窗口,既计算当前批内的“正增量”,也扣除滑出窗口的“负增量”,这样每个时刻的数值都代表最近 60 秒的完整累计。下面这段代码就是 PV + 热点榜的核心:
val windowed = logs .map { case (newsId, _, _, _) => (newsId, 1L) } .reduceByKeyAndWindow( (a: Long, b: Long) => a + b, (a: Long, b: Long) => a - b, Minutes(1), Seconds(5) ) windowed.foreachRDD { rdd => rdd.foreachPartition { it => val jedis = RedisPool.getJedis() it.foreach { case (newsId, cnt) => jedis.zadd("hot_news_topn", cnt.toDouble, newsId) } jedis.expire("hot_news_topn", 7200) RedisPool.returnJedis(jedis) } }逻辑说明分三段。第一段 map 把日志变成 (newsId, 1L) 的键值对,1L 用 Long 类型是为了避免 Int 累加溢出。第二段 reduceByKeyAndWindow 接受四个参数:前两个是加法和减法的函数,Minutes(1) 是窗口长度,Seconds(5) 是滑动间隔,它与 batchDuration 保持一致,保证每个批次滑动一次。第三段 foreachRDD 里的 foreachPartition 是写 Redis 的标准位置——按分区遍历,每个分区拿到一个连接,比逐条创建连接高性能得多。zadd 把 news_id 写入有序集合,score 就是窗口内的 PV 值;expire 设 7200 秒,防止 Redis 里堆积不活跃的新闻。
UV 去重不适合用窗口算子做,因为用户会跨窗口重复出现,精确去重需要全局记录。常见做法是每批把 user_id 写入 Redis 的 Set,key 里带上分钟级时间窗口:
logs.foreachRDD { rdd => rdd.foreachPartition { it => val jedis = RedisPool.getJedis() it.foreach { case (newsId, userId, _, ts) => val minuteKey = s"uv_${ts / 60000}" jedis.sadd(minuteKey, userId) jedis.expire(minuteKey, 120) } RedisPool.returnJedis(jedis) } }这一段用 sadd 把 userId 加入一个以分钟编号命名的 Set,分钟编号来自 ts / 60000,也就是毫秒时间戳除以 60000 取整。因为 Set 天然去重,同一个用户在一分钟内多次点击只算一条。expire 设 120 秒是给这个 key 一个生命周期,避免 Redis 里 minuteKey 无限膨胀。这里的“分钟级 UV”和窗口 PV 的口径不完全一致,答辩时可以说 PV 精确到 5 秒窗口,UV 精确到分钟,理由是全量去重成本高,分钟粒度对运营决策足够。
4.4 背压与 checkpoint:这两个机制决定系统扛不扛得住
运行阶段最怕出现“消费速度大于处理速度”,Kafka 里的日志越积越多,最终内存被撑爆。spark.streaming.backpressure.enabled=true 打开后,Spark 会动态评估任务处理能力,自动调节从 Kafka 拉取的速度;maxRatePerPartition 则是硬上限,当高峰期单分区速率超过 2000 条/秒时直接限流。如果你的模拟器 sleep(0.001) 疯狂灌数据,这两个参数就是保护计算进程不倒的底线。
checkpoint 的坑比很多人想象的深。它不只是“把数据存个盘”,而是周期性保存 DStream 的元数据、未处理的批次和 offset 信息。一旦程序故障重启,Spark 会从最近一个 checkpoint 恢复数据和 offset,这就是 5.1 节里“重启不丢数据”能成立的前提。要注意 checkpoint 目录不要放在临时目录,至少用一个稳定的本地路径,有条件就放 HDFS。另外,checkpoint 恢复要求代码逻辑不能有破坏性变更,如果改了算子结构导致序列化不兼容,恢复可能直接失败,所以演示前别临时改窗口长度。
这一节顺带说明一下和 Structured Streaming 的关系:Spark 2.2 之后的 Structured Streaming 把流处理表达成 DataFrame 上的无界表查询,代码更短、天然支持 sink 到 Kafka 和 Redis。如果要用它,核心是 spark.readStream.format("kafka") 和 groupBy(window($"ts", "60 seconds"), $"news_id")。但 Structured Streaming 的事件时间和 watermark 语义比 DStream 抽象,答辩被深问的几率更高。DStream 的算子更老但更直白,这就是我在 2.1 说选 DStream 更稳的原因。
5. 避坑:实时链路里最容易翻车的 5 个问题与排查步骤
这一章是血泪经验汇总。以下 5 个问题是我见过的高频踩坑点,每条都按现象、原因、解决的顺序写,你可以一边跑一边对着查。
5.1 重启后丢数据或重复消费:offset 到底由谁提交
现象:把 Spark 任务 kill 掉再启动,大屏上的 PV 数值从某个值重新开始,中途几分钟的数据全部消失;或者反过来,重启后数据重复统计了一遍,TopN 的分数翻倍。
原因:这两类现象都指向 offset 管理没做对。如果用的是 createDirectStream 却没有手动提交 offset,进程一死,Kafka 侧发现自己消费组没有提交记录,重启后按 auto.offset.reset=latest 直接跳到最新位置,中间的数据自然就丢了。重复消费则通常是 enable.auto.commit=true 和 checkpoint 恢复同时生效,一批数据算完但没来得及提交,恢复后 Spark 又从 checkpoint 里重放。
解决:在 foreachRDD 处理完指标并写完 Redis 后,拿到 offsetRanges 手动提交。代码是这样:
stream.foreachRDD { rdd => // 先算指标、写 Redis,全部成功后再提交 offset val offsetRanges = rdd.asInstanceOf[HasOffsetRanges].offsetRanges stream.asInstanceOf[CanCommitOffsets].commitAsync(offsetRanges) }注意提交的时机必须在“数据处理成功”之后。如果先 commit 再处理,处理中途失败,那批数据就彻底丢了;如果先处理再 commit,失败后会有极短窗口的重复,但比丢数据好得多。配好 checkpoint 目录再结合这步手动提交,演示和答辩里就可以说“系统支持故障恢复”。
5.2 页面刷新太慢,“实时”变成了“准实时”
现象:大屏数字每隔好几秒才跳一次,有时候等半分钟没变化,和“实时”的宣传完全不符。
原因:常见的有三处叠加。第一,batchDuration 设成了 30 秒,微批周期本身就长;第二,reduceByKeyAndWindow 的窗口长度写的是 Minutes(5),5 分钟内所有数据要累计,数值自然很久不变;第三,前端用 setInterval 每 30 秒轮询一次接口。这三处只要有一处,实时性就打了折扣。
解决:把 batchDuration 调回 5 秒,窗口长度设为 Minutes(1),前端改成 WebSocket 推送而不是轮询。后端每算完一个批次,就把 Redis 里的最新值组装成 JSON 推给所有连着的客户端,浏览器收到数据后 setOption 更新图表。这样从日志进入 Kafka 到大屏变化,中间最多差一个窗口加一个网络往返,用户体感就是秒级刷新。这三个参数是联动的,改的时候一起改,别只调一处。
5.3 热点榜数值和模拟日志对不上:窗口语义和时间戳单位
现象:Kafka 消费端统计某条新闻一共有 120 条日志,Redis 里 hot_news_topn 却只有 100,或者某条不该上榜的新闻排到了前面。
原因:第一可能是 reduceByKeyAndWindow 的减法函数写错,只写了加法函数,导致窗口累计从不衰减,时间一长所有新闻 PV 虚高,越早出现的新闻越占便宜;第二是 ts 的单位在 JSON 解析时被转错,比如毫秒戳被当成秒用,分到错误的时间桶里;第三是窗口滑动间隔和 batchDuration 不一致,导致某些批次的数据没有触发窗口输出。
解决:先核对 reduceByKeyAndWindow 的四个参数,减法函数里的 (a - b) 必须对应窗口滑出的元素,写反了数值只会越来越离谱。再检查 ts 解析环节,Spark 端和 Redis 的 key 都用毫秒时间戳,不要混用秒;模拟器里 time.time() * 1000 已经是毫秒,下游别多做一次除法。最后把窗口参数固定成 60 秒窗口、5 秒滑动,前后端展示的时间标签和窗口对齐。这些字段对不上,基本只能靠打印每条日志的 key 来排查,比较费时间,所以一开始就把 ts 定为毫秒并全程不换单位是节省时间的关键。
5.4 本地能跑,一上集群就 OOM 或任务堆积
现象:本地 IDEA 里运行一切正常,提交到集群后跑十几分钟,Executor 内存使用率飙到 90% 以上,出现 ExecutorLostFailure 或 java.lang.OutOfMemoryError。
原因:本地调试时数据量小、模拟器 QPS 只有每秒 100 条,处理器能轻松消费;上集群时如果还是同一个模拟器,但 Spark 任务的并行度和内存参数没适配,加上没开背压,Kafka 消费端会全速拉取,内存里积压的未处理窗口全部堆积。另一个常见原因是 executor 数量太少,比如默认只申请 1 个 executor 2 核,3 个 Kafka 分区挤在同一个 executor 里处理。
解决:提交前检查三件事。一是 spark-submit 的资源参数,常见做法是 --executor-memory 2g --executor-cores 2 --num-executors 2,单机演示 2 个 executor 就够。二是确认 backpressure 仍在 SparkConf 里,保证消费速率跟随处理能力自动调整。三是把模拟器 QPS 控制稳定,不要开多个模拟器进程同时灌。还有一个小技巧:窗口计算前不需要 persist 中间 RDD,默认的 storageLevel 足够,额外的 persist 只会占用内存。集群部署策略不是节点越多越好,先把一个 executor 跑稳再扩展,这是我个人的血泪经验。
5.5 WebSocket 断连导致大屏卡住:心跳与转发超时
现象:大屏刚打开时图表正常刷新,放一两分钟后完全不动,浏览器控制台提示 WebSocket 连接已关闭,刷新页面又能恢复。
原因:WebSocket 是长连接,但经过 Nginx 反向转发或云服务安全组时,空闲连接会被定时回收。以常用的 Nginx 为例,空闲超过 proxy_read_timeout 默认的 60 秒就会被断开;浏览器侧的 WebSocket 如果一直没有消息,也不再检测断连,大屏就静默卡住。这个现象非常典型,和 Spark 本身没关系,纯属网络链路问题。
解决:给 WebSocket 加心跳机制,后端每隔 30 秒发一条 ping 帧,前端收到后回 pong,连接始终保持活跃。如果走 Nginx 转发,把 proxy_read_timeout 配到 3600 秒。前端在 onclose 事件里做自动重连,避免用户手动刷新。这三个动作做完,大屏可以挂着演示几个小时不掉。测试时不用等自然断连,直接停掉后端进程观察前端是否自动重连,能重连说明兜底逻辑可靠。
6. 让大屏真正“实时”起来:WebSocket 推送 ECharts 的最小闭环与自测核验
可视化大屏的收尾工作,核心在“一条数据从 Redis 到图表”的推送逻辑。后端从 Redis 取出 hot_news_topn 和最近一分钟 PV/UV,组装成 JSON 通过 WebSocket 推给前端,前端拿到后直接 setOption 增量更新。关键点是只更新数据部分,不要让图表整个重新渲染,否则刷新会闪、动画也不连贯。下面这段是 ECharts 侧的核心逻辑:
const chart = echarts.init(document.getElementById("main")); const socket = new WebSocket("ws://localhost:8080/ws/overview"); socket.onmessage = (evt) => { const data = JSON.parse(evt.data); chart.setOption({ xAxis: { data: data.timeLabels }, series: [ { name: "PV", data: data.pvList }, { name: "UV", data: data.uvList } ] }); };这里 setOption 不带第二个参数,ECharts 会保留之前配置做增量合并,动画自然衔接。想让热点榜也动起来,后端推送里带上 topNews 数组,前端在同一个 onmessage 里更新第二个图表。自测方法很简单:让模拟器固定只发 5 条新闻,每轮 20 条消息,跑 60 秒后用 redis-cli 执行 zrevrange hot_news_topn 0 4 withscores,核对分数和页面显示是否一致。再把模拟器停掉,看大屏数值是否停住不再跳动——这也是验证链路是否断在 Kafka 或 Spark 侧的快速手段。
我最早做这类项目时吃过一次亏:把结果直接写进 MySQL,前端轮询 MySQL 接口,结果实时性一塌糊涂,答辩时只能反复强调“架构是实时的”。后来把存储换 Redis、推送换 WebSocket、展示用增量 setOption,才真正让大屏像样地滚动起来。建议你拿到毕设源码后,先跑通这条最小链路,再按需要加地图、渠道饼图和主播监控,方向和参数都对了再扩展。希望帮到你。
本文还有配套的精品资源,点击获取