Spark 2.2 实时新闻分析系统:Kafka+Streaming+MySQL 全链路毕设源码
简介本资源是一套基于Spark 2.2构建的新闻网大数据实时分析系统毕业设计源码面向计算机、大数据及相关专业本科生解决新闻数据流采集、清洗、实时统计与可视化分析等典型教学实践问题。压缩包共34个文件含7个Scala核心业务逻辑文件、6个Java工具类与Flume/HBase集成组件如KfkAsyncHbaseEventSerializer、SimpleRowKeyGenerator、10个依赖jar包、2个XML配置及Web前端相关js/html文件整体3.45MB结构清晰模块覆盖数据接入FlumeKafka、计算Spark Streaming、存储HBase与展示静态页面PNG图表。已有237人学习下载提供完整可运行工程含pom.xml、参考步骤.txt及z_pic目录下的效果截图代码经导师指导与多轮调试具备明确的目录分层如src/main/scala、flume_hbase子模块和典型大数据链路实现范式适合课程设计复现、毕设参考及Spark实时开发入门实战。1. 毕业设计用 Spark 2.2 做新闻网实时分析不是跑通 WordCount 就算“实时”而是把新闻流从 Kafka 拉进来、秒级聚合热点词、写进 MySQL 并同步到 Web 页面——这套源码能直接当毕设答辩底稿也能帮你避开 90% 的 Spark Streaming 踩坑现场很多同学拿 Spark 做毕设一上来就spark-submit --master yarn --class Main xxx.jar结果本地跑通了集群一提交就报NoClassDefFoundError或Kafka offset out of range答辩前两天还在查StreamingContext和SparkSession能不能共存。这套基于 Spark 2.2 的新闻网实时分析系统源码不是玩具 Demo它完整走通了「新闻爬虫模拟 → Kafka 消息入队 → Spark Streaming 消费 窗口统计 → 结果落库 Web 层轮询展示」全链路且所有模块都按毕业设计硬性要求做了分层core流处理核心、daoMySQL 封装、webSpring Boot 简易接口、conf可替换的集群配置。它不依赖 CDH 或 EMR纯原生 Spark 2.2 Kafka 0.10.2 MySQL 5.7 组合连pom.xml里每个依赖的 scopecompile/provided都标得清清楚楚——因为毕设答辩老师真会翻你pom看有没有把spark-yarn_2.11打进 jar 包。如果你正卡在“实时”二字上想证明自己真懂流式计算而不是把离线批处理改个foreachRDD就交差或者被导师问“窗口怎么设才不丢数据”“Exactly-Once 怎么保障”问得哑口无言——这份源码就是你缺的那块拼图。2. 从新闻源模拟到 Kafka 入队为什么不用真实爬虫而用NewsSimulator生成带时间戳、地域标签、热度值的结构化 JSON 流2.1 新闻模拟器的设计逻辑用ScheduledExecutorService控制吞吐量而非while(true)死循环压测毕业设计最怕“数据源不可控”。真实爬虫涉及反爬、IP 封禁、页面结构变动答辩时演示崩了老师一句“你这数据哪来的”就能让你重做。本项目用com.example.simulator.NewsSimulator类替代爬虫核心是ScheduledExecutorService每 2 秒生成一条新闻 JSON// src/main/java/com/example/simulator/NewsSimulator.java public class NewsSimulator { private static final String[] CITIES {北京, 上海, 广州, 深圳, 杭州}; private static final String[] TOPICS {AI, 新能源, 芯片, 教育, 医疗}; public static JSONObject generateNews() { JSONObject news new JSONObject(); news.put(id, UUID.randomUUID().toString().substring(0, 8)); news.put(title, String.format(%s发布%s新政引发%s热议, CITIES[new Random().nextInt(CITIES.length)], TOPICS[new Random().nextInt(TOPICS.length)], new String[]{全民关注, 行业震动, 专家解读, 网友热议}[new Random().nextInt(4)])); news.put(content, 正文约300字含关键词 TOPICS[new Random().nextInt(TOPICS.length)]); news.put(publish_time, System.currentTimeMillis()); // 毫秒时间戳供 Spark 按事件时间窗口 news.put(city, CITIES[new Random().nextInt(CITIES.length)]); news.put(hot_score, new Random().nextInt(100) 50); // 热度值 50~149用于排序 return news; } }提示publish_time是毫秒级时间戳不是new Date()字符串。Spark Streaming 的window和slideDuration必须基于事件时间Event Time否则窗口切分全乱——这是答辩高频扣分点。源码里NewsSimulator生成的每条 JSON 都带这个字段后续KafkaProducer发送时作为value的一部分Spark 消费后可直接.select($publish_time.cast(timestamp))转成时间列。2.2 Kafka 生产端配置acksallretries3保消息不丢但必须配linger.ms5防小包堆积模拟器生成的 JSON 由KafkaProducer推送到news-topic主题。关键不在代码多炫酷而在参数是否符合“毕设级可靠性”// src/main/java/com/example/kafka/KafkaProducerUtil.java Properties props new Properties(); props.put(bootstrap.servers, localhost:9092); props.put(key.serializer, org.apache.kafka.common.serialization.StringSerializer); props.put(value.serializer, org.apache.kafka.common.serialization.StringSerializer); props.put(acks, all); // 所有 ISR 副本写成功才返回防 Leader 挂掉丢数据 props.put(retries, 3); // 网络抖动时重试避免单次失败中断流 props.put(linger.ms, 5); // 等 5ms 攒一批再发平衡延迟与吞吐毕设演示用 5ms 刚好 props.put(batch.size, 16384); // 16KB 批大小配合 linger.ms 防小包泛滥 props.put(buffer.memory, 33554432); // 32MB 缓冲区防生产过快阻塞参数说明acksall是必须项答辩时老师会问“如果 Kafka Leader 宕机你的消息会不会丢”答“默认 acks1”直接不及格linger.ms5是玄学调参点设为 0 则每条消息单独发网络开销大Spark 消费端看到的是大量小批次RDD太碎设为 100ms 则演示时明显卡顿等太久。5ms 是实测平衡点既保证演示流畅又让batch.size生效buffer.memory必须大于batch.size * 10否则send()会阻塞模拟器线程卡死——这是新手常翻车的黑匣子。2.3 Kafka Topic 创建脚本分区数设为 3副本因子为 2匹配 Spark Streaming 的并行度源码包里scripts/create_topic.sh直接创建 topic参数不是随便写的#!/bin/bash # scripts/create_topic.sh KAFKA_HOME/opt/kafka_2.11-0.10.2.1 $KAFKA_HOME/bin/kafka-topics.sh --create \ --zookeeper localhost:2181 \ --replication-factor 2 \ --partitions 3 \ --topic news-topic为什么是 3 分区 2 副本Spark Streaming 的KafkaUtils.createDirectStream每个分区对应一个 RDD partition分区数并行度。3 分区意味着消费端最多启动 3 个 task 并行处理既不过载单核 CPU又足够演示“并行消费”副本因子 2 保证 ZooKeeper 挂一个节点topic 还能读写——毕设演示环境常因虚拟机内存不足 kill 进程副本是后悔药如果你集群只有 1 台机器毕设常见replication-factor必须 ≤ broker 数否则create topic报错Not enough replicas to assign源码已预设为 2你只需确认server.properties里broker.id0且listenersPLAINTEXT://:9092正确。3. Spark Streaming 实时处理核心窗口聚合不是reduceByKeyAndWindow一把梭而是mapWithState管理热点词生命周期3.1 流处理主类NewsStreamingApp的三层架构Input → Process → Output拒绝上帝类源码中com.example.streaming.NewsStreamingApp是入口但它不做任何业务逻辑只负责组装// src/main/scala/com/example/streaming/NewsStreamingApp.scala object NewsStreamingApp extends App { val sparkConf new SparkConf().setAppName(NewsRealTimeAnalysis) val ssc new StreamingContext(sparkConf, Seconds(5)) // batch interval 5s // Input: Kafka Direct Stream val kafkaStream KafkaUtils.createDirectStream[String, String]( ssc, PreferConsistent, Subscribe[String, String](Array(news-topic), kafkaParams) ) // Process: 聚合逻辑封装在 NewsProcessor val resultStream NewsProcessor.process(kafkaStream) // Output: 写 MySQL 更新 Redis 缓存 resultStream.foreachRDD { rdd if (!rdd.isEmpty()) { MySQLDao.saveHotWords(rdd) RedisCache.updateTop10(rdd) } } ssc.start() ssc.awaitTermination() }关键设计点batchDuration5s是刻意设的——比 Kafkalinger.ms5略大确保每个 batch 至少攒够 1~2 条消息避免空 RDD 频繁触发foreachRDDNewsProcessor.process()是独立类把“解析 JSON → 提取标题关键词 → 按城市主题窗口统计”全包进去方便单元测试和答辩时讲清模块职责foreachRDD里加if (!rdd.isEmpty())判空否则空 batch 也会执行saveHotWords导致 MySQL 写入空记录——这是答辩演示时数据表突然多出 10 条NULL的根源。3.2 热点词统计的两种窗口策略滑动窗口保实时性mapWithState保状态一致性源码提供两套统计方案NewsProcessor默认用mapWithState但注释里留了reduceByKeyAndWindow的备选// src/main/scala/com/example/processor/NewsProcessor.scala def process(stream: InputDStream[ConsumerRecord[String, String]]): DStream[(String, Int)] { val parsed stream.map { record val json JSON.parseObject(record.value()) val city json.getString(city) val topic json.getString(topic) // 注意源码实际从 title 提取此处简化示意 (s$city-$topic, 1) } // 方案1mapWithState —— 推荐状态保存在 checkpoint 目录重启不丢计数 val stateSpec StateSpec.function( (key: String, value: Option[Int], state: State[Int]) { val sum value.getOrElse(0) state.getOption.getOrElse(0) state.update(sum) Some((key, sum)) } ).numPartitions(3).timeoutMinutes(10) parsed.mapWithState(stateSpec) // 方案2注释掉reduceByKeyAndWindow —— 简单但状态不持久 // parsed.reduceByKeyAndWindow(_ _, _ - _, Minutes(10), Seconds(30)) }为什么mapWithState是毕设优选reduceByKeyAndWindow的窗口状态存在内存Spark Streaming 重启后归零答辩老师问“系统挂了 10 分钟恢复后热点词排名还准吗”你答“不准”就危险mapWithState的状态序列化到 HDFS 或本地checkpoint目录源码配在conf/spark-streaming.conf重启后自动加载timeoutMinutes(10)表示 key 10 分钟没更新就自动过期防内存泄漏numPartitions(3)强制设为 3匹配 Kafka 的 3 分区避免 shuffle——这是性能关键不设的话默认parallelism200小数据量反而慢。3.3 关键词提取的轻量级实现不用 jieba 或 HanLP用正则 词典双保险新闻标题关键词提取没上 NLP 大模型而是用com.example.util.KeywordExtractor// src/main/java/com/example/util/KeywordExtractor.java public class KeywordExtractor { private static final SetString STOP_WORDS new HashSet(Arrays.asList(发布, 引发, 热议, 新政)); private static final String TOPIC_REGEX (AI|新能源|芯片|教育|医疗); // 预定义主题词典 public static ListString extract(String title) { ListString keywords new ArrayList(); // Step1正则匹配预定义主题词 Matcher m Pattern.compile(TOPIC_REGEX).matcher(title); while (m.find()) { keywords.add(m.group()); } // Step2按顿号、逗号分割过滤停用词 String[] parts title.split([、]); for (String part : parts) { part part.trim(); if (!part.isEmpty() !STOP_WORDS.contains(part)) { keywords.add(part); } } return keywords.stream().distinct().limit(3).collect(Collectors.toList()); } }设计理由毕设不考 NLP 深度考的是“能否在资源受限下合理选型”。jieba 需 Python 环境HanLP 依赖大词典而正则词典 100 行代码搞定mvn package一键打包limit(3)控制每条新闻最多提 3 个词防flatMap后数据爆炸——Spark Streaming 最怕map后flatMap出 100 倍数据executor memory不够直接 OOM停用词列表STOP_WORDS可在conf/stopwords.txt动态维护答辩时老师问“怎么扩展新词”你打开文件现场加一行就行。4. 结果持久化与 Web 展示MySQL 写入不是foreach插入而是JDBCWriter批量 upsert4.1 MySQL DAO 层用PreparedStatement批量 upsert避免 5000 条 insert 变 5000 次网络往返com.example.dao.MySQLDao不用 Hibernate 或 MyBatis纯 JDBC 手写因为毕设要体现“底层可控”// src/main/java/com/example/dao/MySQLDao.java public class MySQLDao { private static final String UPSERT_SQL INSERT INTO hot_words (city_topic, count, last_update) VALUES (?, ?, ?) ON DUPLICATE KEY UPDATE count count VALUES(count), last_update VALUES(last_update); public static void saveHotWords(JavaPairRDDString, Integer rdd) { rdd.foreachPartition(partition - { Connection conn null; PreparedStatement ps null; try { conn DriverManager.getConnection( jdbc:mysql://localhost:3306/news_db?useSSLfalse, root, 123456); ps conn.prepareStatement(UPSERT_SQL); // 批量执行 for (Tuple2String, Integer t : partition) { ps.setString(1, t._1()); ps.setInt(2, t._2()); ps.setLong(3, System.currentTimeMillis()); ps.addBatch(); // 关键攒批 } ps.executeBatch(); // 一次提交 } finally { if (ps ! null) ps.close(); if (conn ! null) conn.close(); } }); } }参数与避坑点ON DUPLICATE KEY UPDATE是核心hot_words表主键是city_topic重复插入自动累加count避免SELECT INSERT/UPDATE的并发冲突ps.addBatch()ps.executeBatch()是性能命脉。如果写成ps.executeUpdate()循环1000 条数据就是 1000 次 TCP 握手演示时卡成 PPTuseSSLfalse必须显式加MySQL 5.7 默认 require SSL不加这参数getConnection直接抛SQLException——这是源码包能跑通但你本地跑不通的头号原因。4.2 Spring Boot Web 层REST 接口返回 JSON不渲染页面用curl或 Postman 即可验证com.example.web.HotWordController只提供两个接口极简// src/main/java/com/example/web/HotWordController.java RestController RequestMapping(/api) public class HotWordController { GetMapping(/top10) public ResponseEntityListHotWord getTop10() { ListHotWord list MySQLDao.queryTop10(); return ResponseEntity.ok(list); } GetMapping(/trend/{city}) public ResponseEntityListHotWord getTrendByCity(PathVariable String city) { ListHotWord list MySQLDao.queryByCity(city); return ResponseEntity.ok(list); } }部署要点application.properties里server.port8081避开了 Spark History Server 的 18080 和 Kafka 的 9092HotWord实体类用DataLombokmvn clean package前确保 IDEA 安装 Lombok 插件否则编译报错——这是新手 IDE 环境没配好的典型翻车接口返回ListHotWord前端用fetch(/api/top10).then(r r.json())即可源码包web/static/index.html里已写好轮询 JSsetInterval(() loadTop10(), 3000)3 秒刷新一次演示时效果直观。4.3 数据库建表 SQLcity_topic设为主键count加索引last_update用TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMPscripts/init_mysql.sql脚本建表字段设计直指痛点-- scripts/init_mysql.sql CREATE DATABASE IF NOT EXISTS news_db CHARACTER SET utf8mb4 COLLATE utf8mb4_unicode_ci; USE news_db; CREATE TABLE hot_words ( city_topic VARCHAR(100) PRIMARY KEY, -- 如 北京-AI count INT NOT NULL DEFAULT 0, last_update TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, INDEX idx_count (count) -- 按 count 排序时走索引 ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;为什么这样建city_topic作主键天然去重INSERT ... ON DUPLICATE KEY UPDATE才生效INDEX idx_count (count)是给ORDER BY count DESC LIMIT 10用的。没这个索引queryTop10()会全表扫描10 万行数据查 10 条要 2 秒演示时卡顿utf8mb4支持 emoji 和生僻字新闻标题可能含“ SpaceX”或“量子”等字符utf8会乱码——答辩老师看到乱码直接质疑数据质量。5. 避坑指南Spark 2.2 Kafka 0.10.2 组合下90% 的失败都源于这 5 个具体错误5.1 现象ClassNotFoundException: org.apache.spark.streaming.kafka010.KafkaUtils原因spark-streaming-kafka-0-10_2.11依赖 scope 写成了compile导致打包时把 Kafka client 打进 fat jar与集群 Spark 自带的 Kafka 版本冲突。解决检查pom.xml确保该依赖 scope 为provideddependency groupIdorg.apache.spark/groupId artifactIdspark-streaming-kafka-0-10_2.11/artifactId version2.2.0/version scopeprovided/scope !-- 必须是 provided -- /dependency注意provided表示运行时由集群提供mvn package不打进 jar。本地测试用mvn compile exec:java集群提交用spark-submit --jars /path/to/spark-streaming-kafka-0-10_2.11-2.2.0.jar显式指定。5.2 现象Spark Streaming 启动后无日志ssc.awaitTermination()一直阻塞原因spark.streaming.kafka.consumer.poll.ms默认 512ms但 Kafka broker 地址配错如localhost:9092在集群模式下应为kafka-host:9092consumer 无法连接poll()永久超时。解决在conf/spark-streaming.conf中显式设置超时并打印 debug 日志spark.streaming.kafka.consumer.poll.ms1000 spark.logConftrue spark.eventLog.enabledtrue然后spark-submit加--conf spark.driver.extraJavaOptions-Dlog4j.debug查看 Kafka 连接日志。5.3 现象MySQL 表里count字段突增 100 倍数据严重失真原因NewsSimulator生成速度过快如ScheduledExecutorService周期设为 100ms而 Spark batch interval 为 5s导致单个 batch 拉取到 50 条消息mapWithState对同一city-topickey 累加 50 次但业务逻辑本意是“每条新闻贡献 1 热度”。解决严格控制模拟器吞吐量。NewsSimulator.java中scheduleAtFixedRate的 delay 设为20002 秒确保 5s batch 内最多 2~3 条新闻。演示前用kafka-console-consumer.sh手动验证消息速率。5.4 现象Web 页面top10接口返回空数组但 MySQL 表里有数据原因MySQLDao.queryTop10()方法里PreparedStatement的ORDER BY count DESC LIMIT 10没加rs.next()判空ResultSet 为空时while(rs.next())不执行list 返回空。解决源码中已修复但若你修改过 DAO请确保ListHotWord list new ArrayList(); while (rs.next()) { // 必须有 rs.next()不能直接 rs.getString() HotWord hw new HotWord(); hw.setCityTopic(rs.getString(city_topic)); hw.setCount(rs.getInt(count)); list.add(hw); } return list;5.5 现象spark-submit报java.lang.NoClassDefFoundError: scala/Product原因Scala 版本不匹配。Spark 2.2 编译于 Scala 2.11但你的pom.xml里scala.version设为 2.12。解决统一 Scala 版本。pom.xml顶部properties中scala.version2.11.12/scala.version spark.version2.2.0/spark.version且所有 Spark 依赖的 artifactId 后缀必须是_2.11如spark-core_2.11、spark-sql_2.11。6. 毕设答辩加分技巧用spark-shell实时验证流处理逻辑三步定位 Kafka 消费断点6.1 第一步用spark-shell直连 Kafka验证消息是否真的在流里别等NewsStreamingApp启动失败才排查先用 Spark Shell 快速验货$SPARK_HOME/bin/spark-shell \ --packages org.apache.spark:spark-sql_2.11:2.2.0,org.apache.spark:spark-streaming-kafka-0-10_2.11:2.2.0 \ --master local[2]进入 shell 后执行import org.apache.spark.streaming._ import org.apache.spark.streaming.kafka010._ val ssc new StreamingContext(sc, Seconds(5)) val kafkaStream KafkaUtils.createDirectStream[String, String]( ssc, PreferConsistent, Subscribe[String, String](Array(news-topic), Map( bootstrap.servers - localhost:9092, group.id - debug-group, key.deserializer - org.apache.kafka.common.serialization.StringDeserializer, value.deserializer - org.apache.kafka.common.serialization.StringDeserializer )) ) kafkaStream.map(_.value()).take(5).foreach(println) ssc.start() ssc.awaitTerminationOrTimeout(10000) // 10秒后自动停作用如果这 5 行能打印出 JSON证明 Kafka 消息正常如果卡住或报错问题在 Kafka 配置不是 Spark 代码。这招能在答辩前 1 小时快速锁定故障域。6.2 第二步用kafka-console-consumer查看 offset确认 consumer group 是否滞留Spark Streaming 的group.id默认是随机的但调试时需固定# 查看 news-topic 的所有 consumer group $KAFKA_HOME/bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --list # 查看指定 group 的 offset假设 group.iddebug-group $KAFKA_HOME/bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \ --group debug-group --describe --topic news-topic关键指标CURRENT-OFFSET当前消费到的位置LOG-END-OFFSET最新消息位置LAG差值即积压消息数。如果LAG 0且持续增长说明 Spark 消费太慢需调大 executor cores 或减小 batchDuration。6.3 第三步在NewsProcessor中加printlncheckpoint目录监控确认状态是否真正持久化mapWithState的状态存在checkpoint目录但默认路径是/tmp/checkpoint重启后可能被清空。源码中已配为hdfs://localhost:9000/checkpoint或本地file:///opt/spark-checkpoint你必须确认conf/spark-streaming.conf中spark.streaming.checkpointDirectory路径可写启动NewsStreamingApp后立刻ls -l /opt/spark-checkpoint应看到00000000000000000000等数字目录手动kill -9进程再启动观察日志是否有Loading state from checkpoint字样。表格checkpoint 目录关键文件含义文件名作用答辩时可展示offsets/存 Kafka 每个 partition 的消费 offset保证 Exactly-Oncecat offsets/00000000000000000000显示 offset 值state/mapWithState的状态快照二进制格式du -sh state/证明状态非空streaming/StreamingContext 元数据ls streaming/证明 checkpoint 目录被识别从那以后我每次重构NewsProcessor都强制走一遍spark-shell验证流、kafka-consumer-groups查 offset、ls checkpoint看状态三步曲——不是为了炫技是避免答辩前夜发现mapWithState根本没生效只能重写逻辑。这套源码的真正价值不在于它多完美而在于它把 Spark Streaming 从黑匣子拆成了可触摸、可验证、可答辩的零件。希望帮到你。本文还有配套的精品资源点击获取