news 2026/9/30 2:55:28

基于Spark2的新闻浏览日志实时分析与可视化系统实战

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
基于Spark2的新闻浏览日志实时分析与可视化系统实战

简介:这份资源是面向大数据方向毕业设计与入门实战的完整项目源码包,围绕新闻网站用户浏览日志,构建从采集、实时流处理到离线分析与可视化的全链路方案。项目以Flume将日志实时写入HBase,再由Spark Streaming消费Kafka或HBase数据流,在Spark 2.x集群上完成实时统计,包括前20名高流量新闻话题、已曝光话题数量及各时段用户浏览量峰值等指标,并借助Spark SQL与Hive完成离线批量计算与历史报表,结果可对接Grafana或前端页面展示。压缩包共35个文件,约3.46MB,以jar依赖、Scala与Java源码为主,辅以xml配置、js与html页面、png可视化图片及md说明文档,目录按flume_hbase、sparkStu、weblogs、z_pic等模块划分,结构清晰。已有68人学习,适合需要完整赛题方案、可运行代码与部署参考步骤的读者,便于快速理解实时与离线分析的数据流组织方式。

1. 新闻浏览日志实时分析:从 Spark2 到可视化大屏的完整落地路径

新闻资讯类产品的后台每天都会沉淀大量浏览日志——谁在什么时间看了哪条新闻、停留多久、从哪个频道点进去、用的什么设备。这些日志单条看没什么价值,但聚合成实时指标之后,就能回答很多运营和产品关心的问题:当前五分钟哪条新闻正在爆、哪个频道的跳出率突然升高、移动端和 PC 端的阅读偏好差多少。这套「基于 Spark2 的新闻浏览日志大数据实时分析与可视化系统」要解决的,就是把原始日志变成可刷新的图表这一整条链路。它适合正在做大数据方向毕业设计的学生,也适合刚接触 Spark 流处理、想找一个完整项目把采集、计算、存储、展示串起来的初中级开发。整条链路的核心技术栈是 Spark2 的 Structured Streaming 或 DStream、Kafka 做缓冲、MySQL 或 HBase 存结果、ECharts 或 Flask 做前端展示。下面按「数据怎么流、代码怎么写、参数怎么调、坑在哪」的顺序拆开讲。

2. 系统分层与数据流:日志从产生到上屏经过哪几层

2.1 四层架构的职责划分

大数据架构通常被拆成采集层、计算层、存储层、展示层四个层次,这套新闻日志系统也不例外。采集层负责把 Nginx 或应用埋点产生的日志收集起来,常见做法是用 Flume 监控日志文件增量,或者用 Logstash 做轻量采集,再统一投递到 Kafka 的一个 topic 里。计算层是 Spark2 的主场,它从 Kafka 消费数据,做窗口聚合、去重、指标计算,把明细日志变成「每分钟各频道 PV」「每五分钟 Top10 新闻」这类结构化结果。存储层承接计算结果,实时性要求高的指标写 MySQL 供前端轮询,数据量大的明细可以落 HBase 或 HDFS。展示层用 Flask 或 SpringBoot 暴露查询接口,前端用 ECharts 画折线图、柱状图、词云。

这四层里最容易出问题的是计算层和存储层的衔接。Spark 算完的结果如果直接写 MySQL,高频写入会打满连接池;如果先攒一批再写,实时性又会下降。我一般会在 Spark 里用foreachPartition批量写入,或者把结果先写 Kafka 再由独立消费者落库,把写入压力从计算任务里剥离出去。

2.2 数据在 Kafka 与 Spark 之间的流转

Kafka 在这套系统里扮演缓冲和削峰的角色。日志产生速率是不均匀的,新闻推送或热点事件发生时可能瞬间暴涨,Spark 任务如果直接对接日志文件,很容易被突发流量打挂。Kafka 把生产者和消费者解耦之后,Spark 可以按自己的节奏消费,积压的数据留在 topic 里不会丢。

一个典型的 topic 设计是:原始日志一个 topic(比如news_log_raw),分区数设成 Spark 消费并行度的整数倍,通常 3 到 6 个分区起步。Spark 的direct模式(Kafka 0.10 之后推荐)会直接读取分区 offset,不经过 ZooKeeper,配合 checkpoint 机制可以实现断点续传。下面这段是 Structured Streaming 从 Kafka 读取并做基础解析的骨架:

# Structured Streaming 读取 Kafka 新闻日志并解析 spark = SparkSession.builder \ .appName("NewsLogRealtime") \ .master("local[4]") \ .getOrCreate() # 从 Kafka 订阅原始日志 topic raw_df = spark.readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", "localhost:9092") \ .option("subscribe", "news_log_raw") \ .option("startingOffsets", "latest") \ .load() # Kafka 的 value 是二进制,转成字符串后按逗号切分字段 parsed_df = raw_df.selectExpr("CAST(value AS STRING) AS log_line") \ .select( split(col("log_line"), ",").getItem(0).alias("user_id"), split(col("log_line"), ",").getItem(1).alias("news_id"), split(col("log_line"), ",").getItem(2).alias("channel"), split(col("log_line"), ",").getItem(3).cast("long").alias("ts") )

这段代码里startingOffsets设成latest表示只消费启动之后的新数据,做实时分析时通常这么设;如果要补历史数据就改成earliest。master("local[4]")是本地调试用的,真正部署到集群要换成yarn并去掉 master 配置。字段切分用split加getItem是最直接的方式,但生产环境日志格式往往带嵌套 JSON,那就得换成from_json配合 schema 定义,否则字段错位会让你排查到怀疑人生。

2.3 窗口聚合与指标计算

日志解析完只是第一步,真正有价值的是聚合。新闻浏览场景最常用的两个窗口是滚动窗口和滑动窗口。滚动窗口(tumbling window)不重叠,适合算「每分钟 PV」;滑动窗口(sliding window)有重叠,适合算「最近 5 分钟 Top10」这种需要连续观察的指标。

# 按 1 分钟滚动窗口统计各频道 PV from pyspark.sql.functions import window, count channel_pv = parsed_df \ .withWatermark("ts", "2 minutes") \ .groupBy( window(col("ts"), "1 minute"), col("channel") ) \ .agg(count("news_id").alias("pv")) # 输出到 MySQL,使用 update 模式让结果持续刷新 query = channel_pv.writeStream \ .outputMode("update") \ .foreachBatch(write_to_mysql) \ .option("checkpointLocation", "/tmp/checkpoint/news_pv") \ .trigger(processingTime="30 seconds") \ .start()

withWatermark("ts", "2 minutes")是处理乱序数据的关键,它告诉 Spark 可以容忍最多 2 分钟的延迟数据,超过这个时间才认为窗口关闭。outputMode("update")只输出有变化的行,比complete模式省资源。trigger设成 30 秒表示每 30 秒触发一次计算,这个值要和窗口长度配合——窗口 1 分钟、触发 30 秒,意味着每个窗口会被计算两次,结果表里会有中间态,前端查询时要注意去重或取最新值。checkpoint 目录必须设,否则任务重启后 offset 丢失会重复消费。

3. 环境搭建与核心代码:把 Spark2 流处理任务跑起来

3.1 集群与依赖版本怎么选

Spark2 这个版本号本身就限定了不少东西。Spark 2.x 最后几个版本是 2.4.x,配套的 Scala 是 2.11 或 2.12,Kafka 客户端建议用 0.10 以上以支持 direct 模式,Hadoop 用 2.7 或 2.8 都比较稳。JDK 必须是 8,Spark2 对 JDK 11 支持不完整,用 JDK 11 跑经常报模块访问错误。Python 侧如果用 PySpark,Python 版本控制在 3.6 到 3.7,再高会和 Spark2 的序列化机制冲突。

依赖这块最容易翻车的是 Kafka 和 Spark 的版本匹配。spark-sql-kafka-0-10这个包是 Structured Streaming 对接 Kafka 用的,版本号要和你 Spark 版本一致,比如 Spark 2.4.5 就配spark-sql-kafka-0-10_2.11:2.4.5。提交任务时用--packages自动拉取,或者提前下好 jar 放到jars目录。

# 提交 Spark2 流处理任务到 YARN 集群 spark-submit \ --master yarn \ --deploy-mode cluster \ --class com.news.NewsLogStreaming \ --executor-memory 2g \ --num-executors 4 \ --executor-cores 2 \ --packages org.apache.spark:spark-sql-kafka-0-10_2.11:2.4.5 \ --files /opt/conf/news.properties \ news-log-analysis.jar

--executor-memory 2g对日志聚合这种轻计算任务够用,如果要做复杂状态计算再往上加。--num-executors 4配合 Kafka 6 个分区,会有 2 个分区排队,实际并行度受 executor 数限制,所以分区数和 executor 数最好成比例。--files把配置文件分发到每个 executor,代码里用相对路径读取,避免硬编码。

3.2 从日志解析到结果落库的完整链路

把前面几段拼起来,一个完整的处理链路是:Kafka 读取 → 字段解析 → 窗口聚合 → 结果写 MySQL。写 MySQL 这一步用foreachBatch比较灵活,可以在每个批次里做批量插入和更新。

# 每个批次把聚合结果批量写入 MySQL def write_to_mysql(batch_df, batch_id): # 转成 Pandas 后用 JDBC 批量插入,减少连接开销 rows = batch_df.collect() if not rows: return conn = pymysql.connect(host="localhost", user="root", password="123456", db="news_analysis") cursor = conn.cursor() sql = """INSERT INTO channel_pv (window_start, window_end, channel, pv) VALUES (%s, %s, %s, %s) ON DUPLICATE KEY UPDATE pv = VALUES(pv)""" for row in rows: cursor.execute(sql, (row["window"]["start"], row["window"]["end"], row["channel"], row["pv"])) conn.commit() cursor.close() conn.close()

ON DUPLICATE KEY UPDATE保证同一个窗口重复计算时更新而不是插入新行,这解决了前面提到的 update 模式重复触发问题。表上要把window_start和channel建联合唯一索引,否则去重不生效。批量写入时如果数据量大,collect()会把所有数据拉到 driver,内存吃紧,更稳的做法是用batch_df.write.jdbc配合mode("append"),但那样就没法做 upsert,需要根据数据量权衡。

3.3 可视化接口与前端刷新

后端用 Flask 暴露一个查询接口,前端定时拉取最新指标。接口逻辑很简单:查 MySQL 最近 N 条记录,按时间排序返回 JSON。

# Flask 查询接口,返回最近 10 分钟的频道 PV @app.route("/api/channel_pv") def channel_pv(): conn = pymysql.connect(host="localhost", user="root", password="123456", db="news_analysis") cursor = conn.cursor(pymysql.cursors.DictCursor) cursor.execute(""" SELECT channel, window_start, pv FROM channel_pv WHERE window_start >= DATE_SUB(NOW(), INTERVAL 10 MINUTE) ORDER BY window_start DESC """) data = cursor.fetchall() cursor.close() conn.close() return jsonify(data)

前端 ECharts 用setInterval每 30 秒请求一次接口,把返回数据映射成折线图的 series。这里有个细节:window_start是 UTC 时间还是本地时间取决于 Spark 的时区配置,如果前端显示时间对不上,先检查spark.sql.session.timeZone这个参数,默认是 UTC,设成Asia/Shanghai才能和本地时间对齐。这个坑我在三个项目里都遇到过,每次都要愣一下才想起来。

4. 避坑与排查:Spark2 流处理最容易翻车的五个地方

4.1 任务重启后数据重复或丢失

现象:Spark 任务因为集群抖动重启后,MySQL 里出现重复的窗口记录,或者某段时间的数据完全缺失。

原因:checkpoint 目录没有配置,或者配置了但被手动删除。Spark 的 offset 提交依赖 checkpoint,没有它任务重启后要么从头消费(重复),要么从 latest 开始(丢失)。

解决:writeStream必须带option("checkpointLocation", "hdfs:///checkpoint/xxx"),路径放在 HDFS 上而不是本地磁盘,否则 executor 换了机器 checkpoint 就找不到了。另外 checkpoint 目录不要和输出目录混用,每个 query 独立一个子目录。

4.2 窗口结果迟迟不输出

现象:数据一直在进,但 MySQL 里就是没有新记录,日志里也看不到报错。

原因:watermark 设得太长,或者数据里的时间戳字段格式不对导致 watermark 无法推进。比如日志时间戳是字符串"2024-01-01 10:00:00",没有 cast 成 timestamp 类型,Spark 无法识别,watermark 永远停在初始值。

解决:解析阶段就把时间字段cast("timestamp"),watermark 延迟设成窗口长度的 1 到 2 倍即可,不要设成 10 分钟这种夸张的值。用query.lastProgress打印进度,看numInputRows和watermark字段确认数据有没有被处理。

4.3 MySQL 连接数被打满

现象:运行一段时间后 Spark 报Too many connections,MySQL 侧看到大量来自 Spark executor 的连接。

原因:foreachBatch里每个批次都新建连接,批次间隔短的时候连接来不及释放。或者 executor 数量多,每个 executor 都持有连接。

解决:把连接创建移到foreachPartition里,一个分区一个连接;或者用连接池(比如 HikariCP)复用连接。更彻底的做法是 Spark 只负责算,结果写 Kafka,由独立的消费者服务落库,把数据库压力从 Spark 任务里彻底剥离。

4.4 中文乱码

现象:MySQL 里存进去的频道名、新闻标题显示成问号或乱码。

原因:Kafka 消息、Spark 解析、MySQL 连接三处编码不一致。常见的是 MySQL 建表时没指定utf8mb4,或者 JDBC 连接串没加characterEncoding=utf8。

解决:建库建表统一用utf8mb4字符集,JDBC URL 加上?useUnicode=true&characterEncoding=utf8,Kafka 生产者侧确认消息是按 UTF-8 编码发送的。三处对齐之后乱码基本就消失了。

4.5 本地能跑集群报 ClassNotFound

现象:spark-submit提交到 YARN 后报ClassNotFoundException,但本地local模式跑得好好的。

原因:依赖 jar 没有分发到集群。--packages在 cluster 模式下有时不会自动分发到所有节点,或者代码里用了本地路径的配置文件。

解决:把依赖 jar 用--jars显式指定,配置文件用--files分发后用相对路径读取。提交前用--verbose看依赖解析结果,确认所有需要的包都在列表里。

5. 让指标更可信:数据质量校验与实时去重技巧

流处理系统跑起来只是及格线,指标能不能信才是关键。新闻日志里有两类脏数据特别常见:一是爬虫或压测流量混进来,把 PV 刷得虚高;二是同一条日志因为采集端重试被重复投递。前者靠user_id白名单或 UA 过滤,后者要在 Spark 里做去重。

去重最直接的方式是用dropDuplicates,但它对流的支持有限,通常要配合 watermark 使用。更稳的做法是用mapGroupsWithState维护一个用户-新闻的访问状态,在状态里判断是否重复。下面是一个简化版的状态去重逻辑:

# 用 mapGroupsWithState 对同一用户短时间内重复浏览去重 from pyspark.sql.streaming import GroupState, GroupStateTimeout def dedup_by_state(key, values, state): # key 是 (user_id, news_id),values 是该组合的所有记录 if state.exists: last_ts = state.get # 5 分钟内重复访问视为同一次,直接丢弃 if values[0]["ts"] - last_ts < 300: return [] state.update(values[0]["ts"]) state.setTimeoutDuration("10 minutes") return [values[0]] dedup_df = parsed_df \ .groupBy("user_id", "news_id") \ .applyInPandasWithState(dedup_by_state, ...)

这段代码的核心思路是给每个「用户-新闻」组合维护一个最后访问时间,5 分钟内的重复访问直接过滤。setTimeoutDuration设成 10 分钟,超过这个时间状态自动清理,避免内存无限增长。实际项目里状态数据量可能很大,要配合 RocksDB 状态存储和合理的超时时间,否则 driver 内存会被状态撑爆。

数据质量校验可以加一个旁路统计:每批次记录总条数、去重后条数、被过滤条数,写入一张监控表。当过滤比例突然升高时,说明采集端可能出了问题,这个信号比指标本身更有价值。我一般会在可视化大屏上留一个小角落放这些质量指标,运营看的是 PV 曲线,开发看的是过滤率曲线,各取所需。

最后说一个我踩过的坑:Spark2 的 Structured Streaming 在update模式下,如果下游是 MySQL 这种不支持事务性 upsert 的存储,窗口结果会出现中间态。我的习惯是给结果表加一个batch_id字段,前端查询时只取每个窗口最大的batch_id,这样即使中间态写进去了,展示出来的也是最终值。这个习惯帮我省了很多解释「为什么 PV 会跳一下又降回去」的口水。希望帮到你。

本文还有配套的精品资源,点击获取

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

Spark+Kafka+Redis实时新闻热点分析系统架构与实战

简介&#xff1a;这是一份基于Apache Spark框架的新闻网大数据实时分析可视化系统项目&#xff0c;面向大数据方向毕业设计、课程设计及推荐算法学习者。项目完整演示了从日志采集、实时流处理到可视化展示的全流程&#xff0c;覆盖Spark Streaming微批处理、Spark SQL数据清洗…

作者头像 李华
网站建设 2026/9/30 2:54:03

神经对话生成对抗性学习复现:从策略梯度到工程落地的完整指南

简介&#xff1a;这是一份机器学习课程设计与期末大作业的高分项目&#xff0c;复现了神经对话生成对抗性学习相关论文。面向计算机、人工智能等专业需要完成对话生成、GAN或论文复现类课题的学生&#xff0c;可同时用于期末大作业、课程设计及毕业设计参考。代码以Python编写&…

作者头像 李华
网站建设 2026/9/30 2:53:05

让 Agent 变成一堆互相调用的服务

文章目录前言一、先说清楚&#xff1a;为什么 while 循环是个"单体"二、Durable Execution 给了我答案的一半三、核心洞察&#xff1a;把 Agent 当成"只有一步"的东西3.1 用 DDD 的话来说这套设计3.2 整体架构四、数据模型&#xff1a;五张表&#xff0c;讲…

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

Random_Poem8

静夜无言烛火熄&#xff0c;冷月高悬千家明。细雨绵绵万户寂&#xff0c;狂风猎猎身心冻。檐下观天窥己命&#xff0c;不觉热泪湿裂甲。别时妻儿亲友送&#xff0c;归家已是赤瓷人。翌日清晨妻儿泣&#xff0c;丈夫仰首笑永恒。

作者头像 李华