ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

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

基于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(processingTime30 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.12Kafka 客户端建议用 0.10 以上以支持 direct 模式Hadoop 用 2.7 或 2.8 都比较稳。JDK 必须是 8Spark2 对 JDK 11 支持不完整用 JDK 11 跑经常报模块访问错误。Python 侧如果用 PySparkPython 版本控制在 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(hostlocalhost, userroot, password123456, dbnews_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(hostlocalhost, userroot, password123456, dbnews_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 connectionsMySQL 侧看到大量来自 Spark executor 的连接。原因foreachBatch里每个批次都新建连接批次间隔短的时候连接来不及释放。或者 executor 数量多每个 executor 都持有连接。解决把连接创建移到foreachPartition里一个分区一个连接或者用连接池比如 HikariCP复用连接。更彻底的做法是 Spark 只负责算结果写 Kafka由独立的消费者服务落库把数据库压力从 Spark 任务里彻底剥离。4.4 中文乱码现象MySQL 里存进去的频道名、新闻标题显示成问号或乱码。原因Kafka 消息、Spark 解析、MySQL 连接三处编码不一致。常见的是 MySQL 建表时没指定utf8mb4或者 JDBC 连接串没加characterEncodingutf8。解决建库建表统一用utf8mb4字符集JDBC URL 加上?useUnicodetruecharacterEncodingutf8Kafka 生产者侧确认消息是按 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 会跳一下又降回去」的口水。希望帮到你。本文还有配套的精品资源点击获取
返回列表