
简介基于Spark2.2的新闻网大数据实时分析系统源码面向大数据方向毕业设计学生、Spark入门开发者以及需要实时日志分析方案的技术人员。以新闻网站用户行为数据为切入点覆盖数据采集、存储、处理与热点趋势挖掘全流程可用于毕业设计参考、技术复现或课程实验。压缩包共43个文件、约3.64MB核心为10个JAR运行依赖、7个Scala与6个Java源码文件辅以ZIP附赠资料、备份文件与说明文档其中flume_hbase目录展示Flume与HBase集成实现日志采集与持久化weblogs提供模拟新闻日志sparkStu收录Spark学习实践代码参考步骤与README帮助快速掌握运行逻辑。相较于MapReduceSpark2.2在内存计算与流处理上更具优势该工程正好演示了这种新一代框架的落地方式。项目经导师认可并严格调试代码可运行能帮助深入理解Spark实时计算与HBase存储的协同也便于迁移到其他大数据场景已有54人学习使用适合巩固Spark实践能力并完成设计型任务。1. 基于Spark2.2的新闻网大数据实时分析系统到底在解决什么问题“基于Spark2.2的新闻网大数据实时分析系统”这类项目标题每年大数据毕业设计和中小团队自建报表里都能见到Kafka 收新闻点击日志Spark Streaming 做 PV/UV、分类热度统计结果进 Redis 和 MySQL再由可视化页面展示。做这套的人通常不是要追新特性而是想要一个能在 Spark 2.2 老集群上稳定跑、代码改得动的实时分析链路。系统本身不复杂真正难的是 StreamingContext 配置、窗口聚合的边界参数以及“先入库还是先提交 offset”这类顺序问题。这篇笔记面向要做毕业设计、或从零搭一套实时分析系统的开发者按一套可复现的源码方案讲清楚从埋点数据到实时指标哪些环节必须自己写哪些坑值得提前绕开。2. 先把实时分析链路拆开双链路架构与存储选型2.1 为什么 Spark2.2 适合当这套系统的实时计算层选 Spark 2.2 不是因为它新而是因为它在“实时分析”这个场景里足够成熟。2.2 版本的 Spark StreamingDStream已经非常稳定搭配 Kafka 0.10 的 Direct API能做到 At-Least-Once 语义下的精准恢复而 2.2 里 Structured Streaming 还带 Experimental 标签流式 join、状态管理的坑还没填完。常见做法是用 DStream 做主计算层后续再按需迁移。对新闻报道这类高吞吐、允许轻微重复统计的场景DStream 的生态和资料量反而是优势——搜一个问题几乎都有答案。还有一个现实因素很多实验室和公司的生产集群还停在 Spark 2.x 早期版本。Spark 2.2 用 Scala 2.11 编译和 CDH 5.x、老版本 HDFS 兼容性好部署成本低。这套系统的设计目标是“能跑、能改、能演示”不是“能上 TPC 榜单”所以稳定性优先于性能。后面所有代码都按 Spark 2.2 Scala 2.11 Kafka 0.10 来写。2.2 消息模型新闻点击日志的 Topic 设计与 JSON 格式实时分析的第一步是把数据送进 Kafka。新闻站的埋点日志一般长这样用户点了哪条新闻、属于哪个频道、什么时间点的、用户标识是什么。把前端埋点和服务端 nginx 日志统一成一条 JSON推送进 Kafka 的news_click主题。{user_id:u_10086,news_id:n_20301,category:sports,action:click,ts:1693209600123,source:app}设计时注意几点source用来区分 App 端还是 Web 端后面分流统计要用ts用毫秒时间戳不要用字符串时间减少解析成本category必须由后端保证枚举值统一前端传什么就存什么后面聚合会省掉一堆数据清洗的事。主题分区数建议和 Spark 执行器总数对齐。比如 3 个 Worker、每个 Worker 2 核分区数就设 6。分区数大于消费并发度会浪费少于消费并发度会造成部分 Executor 空闲。这里我一般会先给 6 个分区后面压测再调。2.3 存储选型Redis 扛实时指标MySQL 扛历史报表这套系统里 Redis 和 MySQL 分工明确Redis 负责毫秒级读写的实时指标MySQL 负责落历史报表。Redis 里主要放三类数据——PV 计数器、UV 去重集合、热度榜有序集合。PV 用INCR做原子自增UV 用PFADD或者SADD去重热度榜用ZINCRBY做增量排序三个数据结构正好覆盖实时侧全部需求。MySQL 侧则是一张宽表按时间窗口存聚合结果供报表系统按天、按小时查询。表结构下面给出这是最常被抄走的部分。模块技术选型理由日志采集Flume / 模拟脚本Flume 适合接 nginx 日志模拟脚本适合本地复现消息队列Kafka 0.10Spark 2.2 官方支持最好Direct API 成熟实时计算Spark Streaming DStream2.2 下稳定窗口算子齐全实时存储Redis原子计数、集合去重、有序榜单一站式解决历史存储MySQL报表查询简单团队都会用可视化ECharts前端拉 Redis/MySQL 接口出图和这套链路无关2.4 项目源码目录建议这套系统的源码建议按标准 sbt 工程组织我复现时用的目录结构如下直接抄这个结构能少踩很多依赖坑。路径职责src/main/scala/com/news/streaming/实时计算主程序src/main/scala/com/news/util/Redis、MySQL、Kafka 工具类src/main/resources/spark-defaults.conf、日志配置src/test/scala/本地模式单测build.sbt依赖管理锁定 Spark 2.2.0依赖锁定的意义很大。Spark 2.2 如果直接拉 2.4 的spark-streaming-kafka-0-10会出现二进制不兼容Kafka client 版本最好也和集群一致。下面给出可直接用的build.sbt核心依赖。name : news-realtime-analysis version : 1.0 scalaVersion : 2.11.12 libraryDependencies Seq( org.apache.spark %% spark-core % 2.2.0 % provided, org.apache.spark %% spark-streaming % 2.2.0 % provided, org.apache.spark %% spark-streaming-kafka-0-10 % 2.2.0, redis.clients % jedis % 2.9.0, mysql % mysql-connector-java % 5.1.49 ) assemblyMergeStrategy in assembly : { case PathList(META-INF, xs _*) MergeStrategy.discard case _ MergeStrategy.first }参数说明spark-core和spark-streaming用provided因为集群的 Spark 环境会提供这两个 jarspark-streaming-kafka-0-10必须随应用提交它不在 Spark 安装目录里。Jedis 2.9 对 Redis 3.x 支持最稳MySQL 驱动用 5.1.x不要用新版 8.x——它和 Spark 2.2 的序列化机制偶尔有兼容性问题这是踩过的坑。3. Spark2.2 实时计算核心代码从 Kafka 接入到窗口聚合3.1 构建 StreamingContext批处理间隔决定延迟上限实时计算的第一步是创建StreamingContext。批处理间隔是整个系统的节拍器它决定了延迟上限。新闻点击场景 5 秒一批是性价比最高的选择1 秒一批对 2.2 的调度开销太大10 秒一批做 PV/UV 又不够“实时”。import org.apache.spark.SparkConf import org.apache.spark.streaming.{Seconds, StreamingContext} object NewsRealtimeApp { def main(args: Array[String]): Unit { val conf new SparkConf() .setAppName(NewsRealtimeAnalysis) .setIfMissing(spark.master, local[2]) .set(spark.serializer, org.apache.spark.serializer.KryoSerializer) .set(spark.streaming.backpressure.enabled, true) .set(spark.streaming.kafka.maxRatePerPartition, 2000) val ssc new StreamingContext(conf, Seconds(5)) ssc.checkpoint(hdfs:///tmp/news-streaming-checkpoint) // 后续业务逻辑都挂在 ssc 上 ssc.start() ssc.awaitTermination() } }参数说明setIfMissing(spark.master, local[2])是为了本地调试不显式传 master 也能跑线上提交时用spark-submit --master yarn覆盖。KryoSerializer能省大量内存因为 DStream 每个批次的对象都要序列化到网络或磁盘。maxRatePerPartition是关键限速参数防止 Kafka 积压时瞬间拉爆 Executor。checkpoint目录存 DStream 操作图和 RDD 元数据后面避坑章会详细说它的副作用。3.2 Kafka Direct 模式接入不依赖 ZK 的 offset 管理Spark 2.2 接 Kafka 0.10 要用KafkaUtils.createDirectStreamDirect 模式的好处是按分区拉取不经过 Receiver也就没有 WAL 和 Receiver 故障恢复的问题。offset 由 Spark 自己管理语义更清晰。import org.apache.kafka.common.serialization.StringDeserializer import org.apache.spark.streaming.kafka010._ val kafkaParams Map[String, Object]( bootstrap.servers - node1:9092,node2:9092,node3:9092, key.deserializer - classOf[StringDeserializer], value.deserializer - classOf[StringDeserializer], group.id - news-realtime-group, auto.offset.reset - latest, enable.auto.commit - (false: java.lang.Boolean) ) val topics Array(news_click) val stream KafkaUtils.createDirectStream[String, String]( ssc, LocationStrategies.PreferConsistent, ConsumerStrategies.Subscribe[String, String](topics, kafkaParams) )逻辑说明enable.auto.commit必须设为 false这是整个系统的关键决策。如果让 Kafka 自动提交 offset而 Spark 这边还没处理完一旦 Executor 崩溃就会丢数据关闭自动提交后由 Spark 在每个批次处理完后手动提交才能保证“处理完再提交”的顺序。auto.offset.reset设latest因为实时报表只关心新数据重跑历史用离线链路做。3.3 PV/UV 与热度 TopN 的窗口计算新闻热度不能只看累计点击量还要看“最近一段时间”的表现。这里用reduceByKeyAndWindow做滑动窗口聚合统计 60 秒窗口内每篇新闻的点击数再把结果写进 Redis 的 ZSet按分数排序就能拿到实时热度榜。case class ClickLog(userId: String, newsId: String, category: String, ts: Long) val parsedStream stream.map(record { val json JSON.parseObject(record.value()) ClickLog( json.getString(user_id), json.getString(news_id), json.getString(category), json.getLong(ts) ) }) // 窗口 60 秒滑动 30 秒统计每篇新闻的点击量 val newsClickCount parsedStream .map(log (log.newsId, 1L)) .reduceByKeyAndWindow( (a: Long, b: Long) a b, (a: Long, b: Long) a - b, Seconds(60), Seconds(30) ) // 取每个批次窗口结果中的 TopN写入 Redis newsClickCount.foreachRDD { rdd val topN rdd.sortBy(_._2, ascending false).take(20) // topN 写入 Redis ZSetkey 为 hot_rank:当前窗口 }参数说明reduceByKeyAndWindow有四个参数前两个是窗口内增量和减量函数第三个是窗口长度 60 秒第四个是滑动间隔 30 秒。窗口长度和滑动间隔必须是批处理间隔的整数倍这里批处理 5 秒60 和 30 都能整除不会报错。减法函数的原理是维护“旧窗口移除 新窗口加入”比每个窗口全量重算省一半资源。sortBy这里的 TopN 是近似值——它在每个 Spark 分区内局部排序后再合并但对 60 秒窗口的热度榜够用。3.4 结果写 Redis 和 MySQL幂等写入与批量提交实时统计结果要同时进 Redis 和 MySQL。Redis 侧的核心操作是原子计数和去重MySQL 侧则是按主键做幂等更新。// Redis 写入每个分区一个连接避免每条记录都建连接 parsedStream.foreachRDD { rdd rdd.foreachPartition { partition val jedis RedisPool.getResource partition.foreach { log val pvKey spv:${log.category}:${currentMinute} jedis.incr(pvKey) // 用 Set 做近似 UV每天一个 key避免无限增长 jedis.sadd(suv:${log.newsId}:${today}, log.userId) // 热度榜直接用 ZINCRBY 加权重 jedis.zincrby(hot_rank, 1.0, log.newsId) } jedis.close() } }逻辑说明foreachPartition是这里最重要的写法。如果写成foreach每条点击日志都会从连接池拿一次连接Redis 吞吐直接被打爆按分区处理让每个 Executor 的分区只复用同一个 Jedis 连接连接数等于分区数而不是数据条数。PV的 key 带分类和分钟是按“频道 × 分钟”粒度做实时趋势分析UV的 key 带日期是为了让集合能自然过期不会无限膨胀。ZINCRBY每来一条点击给对应新闻加 1 分Redis 自动维护排序报表端直接ZREVRANGE hot_rank 0 19拿前 20 名。MySQL 侧需要一张结果表建表语句如下按“日期 新闻 窗口”做唯一键写入用REPLACE INTO保证重复执行不会产生脏数据。CREATE TABLE news_window_stats ( id BIGINT PRIMARY KEY AUTO_INCREMENT, stat_date VARCHAR(10) NOT NULL, news_id VARCHAR(32) NOT NULL, category VARCHAR(20) NOT NULL, pv BIGINT DEFAULT 0, uv BIGINT DEFAULT 0, heat_score DOUBLE DEFAULT 0, window_start BIGINT NOT NULL, window_end BIGINT NOT NULL, update_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, UNIQUE KEY uk_date_news_window (stat_date, news_id, window_start, window_end) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;字段说明window_start和window_end存的是窗口起始和结束的毫秒时间戳这样同一个新闻在 10:00:00 和 10:01:00 两个窗口的记录可以共存。UNIQUE KEY加上REPLACE INTO就是幂等写入的保障。报表系统要做小时级趋势直接按stat_date和window_start分组查询就行不需要再对原始日志做二次计算。4. 部署与调参这套系统在集群上的落地参数4.1 大数据集群部署策略三节点怎么摆很多人的集群规划是把所有服务塞进三台机器能跑通但性能很玄学。这里给出一个新闻实时分析系统最常见的三节点部署策略两台中高端配置的机器做计算和存储混部一台低配机器做调度。节点部署组件资源建议node1NameNode、ResourceManager、Kafka Broker16C/32G系统盘 2 块数据盘node2DataNode、NodeManager、Kafka Broker、Redis16C/32Gnode3DataNode、NodeManager、MySQL、ZooKeeper8C/16G这套布局的思路是让 Spark Executor 尽量和 Kafka 在同一批机器上拉取消息走内网带宽。Kafka 本身需要至少 2 个副本放在 node1 和 node2 正好满足replication.factor2。Redis 单独跟 MySQL 放一起因为实时链路的写吞吐主要由 Redis 抗和 MySQL 同节点能减少一次跨机网络耗时。ZooKeeper 在 node3 单点对这套系统不是瓶颈Kafka 的 controller 选举只在 broker 宕机时触发平时不参与数据路径。4.2 spark-submit 与 Streaming 核心参数实时任务用spark-submit提交到 YARN推荐client模式这样日志直接打到终端调试方便生产再用cluster模式。spark-submit \ --class com.news.streaming.NewsRealtimeApp \ --master yarn \ --deploy-mode client \ --driver-memory 2g \ --executor-memory 4g \ --executor-cores 2 \ --num-executors 3 \ --conf spark.streaming.backpressure.enabledtrue \ --conf spark.streaming.kafka.maxRatePerPartition2000 \ --conf spark.streaming.receiver.maxRate2000 \ --conf spark.serializerorg.apache.spark.serializer.KryoSerializer \ --conf spark.streaming.kafka.consumer.cache.initialCapacity16 \ news-realtime-analysis.jar每个参数都有具体作用executor-cores2不要设太大因为每个 Executor 的 CPU 还要留给 GC 和 Kafka 拉取线程num-executors3对应 3 个节点和前面 6 个 Kafka 分区刚好形成每个 Executor 消费 2 个分区的映射。maxRatePerPartition2000意味着每个分区每秒最多消费 2000 条3 个 Executor 总共每秒最多消费 12000 条这对新闻站完全够用——如果实际峰值超出优先加分区而不是调大限速因为限速是最终兜底不是性能目标。consumer.cache.initialCapacity是 Kafka consumer 缓存的初始容量太小会在高并发下频繁扩容浪费 GC。4.3 checkpoint 与 offset 的持久化位置Spark Streaming 的 checkpoint 目录同时保存两类东西一是 DStream 操作图二是 RDD 数据和 offset。这意味着 checkpoint 既是后悔药也是紧箍咒。后悔药的部分任务挂了可以从最近一次 checkpoint 恢复不用重新消费整段 Kafka 数据。紧箍咒的部分如果你改了代码里 DStream 的转换逻辑再从这个 checkpoint 恢复Spark 会发现操作图和当前代码不一致直接抛异常解决方式是删掉 checkpoint 目录重新跑。这放在生产里是没法接受的所以我的做法是——checkpoint 只做 RDD 数据恢复offset 单独持久化到外部存储。// 每个批次处理完成后手动把 offset 写入 Redis stream.foreachRDD { rdd val offsetRanges rdd.asInstanceOf[HasOffsetRanges].offsetRanges // 业务写入Redis 指标、MySQL 报表 // 业务写完再提交 offset stream.asInstanceOf[CanCommitOffsets].commitAsync(offsetRanges) }这段代码的顺序很关键必须先完成 Redis 和 MySQL 的写入再调用commitAsync提交 offset。如果先提交后写入写入失败的那批数据就永久丢了如果先写入后提交写入成功但提交失败会重复消费但重复消费大多数场景可以容忍丢失不可容忍。写 Redis 的 offset 键可以设计成offset:{topic}:{partition}value 存当前消费到的 offset重启时从这里恢复。5. 实时链路避坑五个让我半夜爬起来看日志的问题5.1 重启后重复消费计数被翻倍现象任务重启后Redis 里的 PV 每分钟涨得比以前快数值接近翻倍。原因enable.auto.commit之前设成了 trueKafka 自动提交 offset 的时机和 Spark 真正处理完的时机不一致。Spark 拉取消息、Executor 还在计算时Kafka 已经按 poll 周期提交了 offset任务重启后Kafka 认为这些消息已经消费过但实际结果没落库于是重新拉取和旧数据叠加。解决把enable.auto.commit设为 false改用commitAsync在批次处理完、业务写入完成后提交。这套系统允许轻微重复统计但绝不允许丢数据所以顺序必须是“写库在前提交 offset 在后”。5.2 窗口时长不是批次间隔的整数倍直接抛异常现象reduceByKeyAndWindow传了窗口 60 秒、滑动 20 秒启动后立刻报IllegalArgumentException: Requirement failed。原因Spark Streaming 要求窗口长度和滑动间隔必须是批处理间隔的整数倍。批处理间隔 5 秒滑动 20 秒能被整除窗口 60 秒也能被整除但如果把批处理间隔改成 7 秒或者窗口改成 65 秒就会触发这个异常。解决改参数前先算一下。窗口和滑动都除以批处理间隔结果必须是正整数。更安全的做法是把批处理间隔、窗口长度、滑动间隔定义成常量放在配置里统一修改val batchInterval Seconds(5) val windowLength Seconds(60) // 60 / 5 12 val slideInterval Seconds(30) // 30 / 5 6这样从源头避免参数不整除造成的启动失败。5.3 checkpoint 恢复后 Kafka 不再消费现象改了业务代码重新提交任务从 checkpoint 恢复运行日志显示 task 在跑但 Kafka 的 consumer group 的 Lag 持续增长没有新数据进入统计。原因checkpoint 里保存的是旧 DStream 操作图代码变更后 Spark 从 checkpoint 恢复操作图时发现逻辑对不上会以 checkpoint 里的旧逻辑为准在某些版本下表现为 consumer 初始化异常或 offset 读取异常进程还活着但不再拉取数据。解决不要依赖 checkpoint 恢复 offset。把 offset 独立持久化到 Redis任务启动时按 Redis 里的 offset 重建 consumercheckpoint 目录只作为 RDD 数据的缓存代码变更时删除 checkpoint 目录。具体操作是任务初始化时读 Redis 里的offset:{topic}:{partition}用它构建ConsumerStrategies.Subscribe的 starting offsets每次批次处理后回写新 offset。这样重启后恢复的是业务状态而不是被固化死的操作图。5.4 输出算子每条记录建一次连接Kafka 积压拉爆内存现象处理延迟指标持续上升Kafka 消费 lag 越来越大查看 Executor 日志发现有大量Too many connections和 GC 告警。原因常见的初学者写法是在foreach里写 Redis 或 MySQL每一条点击日志都new Jedis(...)或DriverManager.getConnection(...)。连接建立开销远大于数据处理本身导致单批次处理时间超过批次间隔积压像滚雪球一样扩大。解决所有外部存储写入必须走foreachPartition每个分区一个连接连接池用完归还不要new一个丢一个。Jedis 连接池初始化放在 Driver 端通过广播变量或单例传给 Executor代码如下rdd.foreachPartition { partition val jedis RedisPool.getResource try { partition.foreach { log /* Redis 写入 */ } } finally { jedis.close() // close 归还连接池不是断开连接 } }同样的逻辑适用于 MySQLforeachPartition内获取一个Connection批量执行REPLACE INTO提交事务后再关闭连接。把连接的复用粒度从“每行”提升到“每分区”写入性能通常有 10 倍以上的提升。5.5 处理延迟像滚雪球一样涨恢复不了现象任务跑了一个小时Processing Time从 2 秒慢慢涨到 30 秒批处理间隔 5 秒每批都积压调大资源后短暂恢复过一阵又涨上去。原因这是典型的背压没生效。Direct 模式下如果 Kafka 分区数和 Executor 并发不匹配或者maxRatePerPartition设得过大拉取速率超过处理能力延迟就会累加。资源加得再多单分区消费速率仍然由 Kafka 分区数决定瓶颈没解决。解决先看瓶颈在哪。如果 CPU 使用率低但延迟高大概率是序列化或外部存储慢——检查是否用了 Kyro检查 Redis 连接池大小如果 CPU 已经打满说明计算资源不够——优先增加 Kafka 分区数让消费并发能力先上来再考虑加 Executor。backpressure.enabled要开着但别期望它自动搞定一切——它调节的只是消费速率阈值真实吞吐上限仍然取决于你的处理逻辑和外部存储性能。定位瓶颈时用Batch Processing Time和Scheduling Delay两个指标对照着看前者高说明处理慢后者高说明资源排队。6. 再往上走一步验证手段与 DStream 向结构化流式的扩展这套系统跑通之后下一步值得做的有两件事一是建立一套可重复的验证流程二是评估是否要从 DStream 迁移到 Structured Streaming。验证流程我一般用 Kafka 自带的生产者脚本灌模拟数据。写一个 Python 或 shell 脚本按真实新闻点击的速率往news_click主题发数据观察 Redis 里的 PV 计数是否符合预期、ZSet 的热度排名是否按窗口滑动更新然后故意停掉 Executor再重启看 MySQL 里有没有重复更新但无丢失。这个过程不需要写测试框架一条命令就能完成kafka-console-producer.sh \ --broker-list node1:9092,node2:9092 \ --topic news_click \ --property parse.keytrue \ --property key.separator:启动后逐行输入 JSON 格式的点击消息5 秒后去 Redis 查对应的pv:{category}:{minute}和hot_rank能立刻确认整条链路的数据流是否完整。我每次上线前都会先用这种方式跑 10 分钟确认 offset 单调递增、Redis 指标无跳变再放真实流量。至于迁移 Structured StreamingSpark 2.2 里它还是 Experimental不建议直接上生产但代码结构可以提前预留抽象层——把parsedStream的类型从DStream[ClickLog]抽象成Stream[ClickLog]将来迁到新版本时窗口聚合逻辑可以整段复用。这套基于 Spark2.2 的新闻网大数据实时分析系统真正值钱的部分不是框架本身而是“写入前提交 offset”“分区内复用连接”“窗口参数按批次对齐”这些经过线上验证的细节它们才是实时链路不翻车的底牌。希望帮到你。本文还有配套的精品资源点击获取