時(shí)新聞分析系統(tǒng):Kafka+Streaming+MySQL 全鏈路畢設(shè)源碼)
簡(jiǎn)介本資源是一套基于Spark 2.2構(gòu)建的新聞網(wǎng)大數(shù)據(jù)實(shí)時(shí)分析系統(tǒng)畢業(yè)設(shè)計(jì)源碼面向計(jì)算機(jī)、大數(shù)據(jù)及相關(guān)專業(yè)本科生解決新聞數(shù)據(jù)流采集、清洗、實(shí)時(shí)統(tǒng)計(jì)與可視化分析等典型教學(xué)實(shí)踐問(wèn)題。壓縮包共34個(gè)文件含7個(gè)Scala核心業(yè)務(wù)邏輯文件、6個(gè)Java工具類與Flume/HBase集成組件如KfkAsyncHbaseEventSerializer、SimpleRowKeyGenerator、10個(gè)依賴jar包、2個(gè)XML配置及Web前端相關(guān)js/html文件整體3.45MB結(jié)構(gòu)清晰模塊覆蓋數(shù)據(jù)接入FlumeKafka、計(jì)算Spark Streaming、存儲(chǔ)HBase與展示靜態(tài)頁(yè)面PNG圖表。已有237人學(xué)習(xí)下載提供完整可運(yùn)行工程含pom.xml、參考步驟.txt及z_pic目錄下的效果截圖代碼經(jīng)導(dǎo)師指導(dǎo)與多輪調(diào)試具備明確的目錄分層如src/main/scala、flume_hbase子模塊和典型大數(shù)據(jù)鏈路實(shí)現(xiàn)范式適合課程設(shè)計(jì)復(fù)現(xiàn)、畢設(shè)參考及Spark實(shí)時(shí)開(kāi)發(fā)入門(mén)實(shí)戰(zhàn)。1. 畢業(yè)設(shè)計(jì)用 Spark 2.2 做新聞網(wǎng)實(shí)時(shí)分析不是跑通 WordCount 就算“實(shí)時(shí)”而是把新聞流從 Kafka 拉進(jìn)來(lái)、秒級(jí)聚合熱點(diǎn)詞、寫(xiě)進(jìn) MySQL 并同步到 Web 頁(yè)面——這套源碼能直接當(dāng)畢設(shè)答辯底稿也能幫你避開(kāi) 90% 的 Spark Streaming 踩坑現(xiàn)場(chǎng)很多同學(xué)拿 Spark 做畢設(shè)一上來(lái)就spark-submit --master yarn --class Main xxx.jar結(jié)果本地跑通了集群一提交就報(bào)NoClassDefFoundError或Kafka offset out of range答辯前兩天還在查StreamingContext和SparkSession能不能共存。這套基于 Spark 2.2 的新聞網(wǎng)實(shí)時(shí)分析系統(tǒng)源碼不是玩具 Demo它完整走通了「新聞爬蟲(chóng)模擬 → Kafka 消息入隊(duì) → Spark Streaming 消費(fèi) 窗口統(tǒng)計(jì) → 結(jié)果落庫(kù) Web 層輪詢展示」全鏈路且所有模塊都按畢業(yè)設(shè)計(jì)硬性要求做了分層core流處理核心、daoMySQL 封裝、webSpring Boot 簡(jiǎn)易接口、conf可替換的集群配置。它不依賴 CDH 或 EMR純?cè)?Spark 2.2 Kafka 0.10.2 MySQL 5.7 組合連pom.xml里每個(gè)依賴的 scopecompile/provided都標(biāo)得清清楚楚——因?yàn)楫呍O(shè)答辯老師真會(huì)翻你pom看有沒(méi)有把spark-yarn_2.11打進(jìn) jar 包。如果你正卡在“實(shí)時(shí)”二字上想證明自己真懂流式計(jì)算而不是把離線批處理改個(gè)foreachRDD就交差或者被導(dǎo)師問(wèn)“窗口怎么設(shè)才不丟數(shù)據(jù)”“Exactly-Once 怎么保障”問(wèn)得啞口無(wú)言——這份源碼就是你缺的那塊拼圖。2. 從新聞源模擬到 Kafka 入隊(duì)為什么不用真實(shí)爬蟲(chóng)而用NewsSimulator生成帶時(shí)間戳、地域標(biāo)簽、熱度值的結(jié)構(gòu)化 JSON 流2.1 新聞模擬器的設(shè)計(jì)邏輯用ScheduledExecutorService控制吞吐量而非while(true)死循環(huán)壓測(cè)畢業(yè)設(shè)計(jì)最怕“數(shù)據(jù)源不可控”。真實(shí)爬蟲(chóng)涉及反爬、IP 封禁、頁(yè)面結(jié)構(gòu)變動(dòng)答辯時(shí)演示崩了老師一句“你這數(shù)據(jù)哪來(lái)的”就能讓你重做。本項(xiàng)目用com.example.simulator.NewsSimulator類替代爬蟲(chóng)核心是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, 新能源, 芯片, 教育, 醫(yī)療}; public static JSONObject generateNews() { JSONObject news new JSONObject(); news.put(id, UUID.randomUUID().toString().substring(0, 8)); news.put(title, String.format(%s發(fā)布%s新政引發(fā)%s熱議, CITIES[new Random().nextInt(CITIES.length)], TOPICS[new Random().nextInt(TOPICS.length)], new String[]{全民關(guān)注, 行業(yè)震動(dòng), 專家解讀, 網(wǎng)友熱議}[new Random().nextInt(4)])); news.put(content, 正文約300字含關(guān)鍵詞 TOPICS[new Random().nextInt(TOPICS.length)]); news.put(publish_time, System.currentTimeMillis()); // 毫秒時(shí)間戳供 Spark 按事件時(shí)間窗口 news.put(city, CITIES[new Random().nextInt(CITIES.length)]); news.put(hot_score, new Random().nextInt(100) 50); // 熱度值 50~149用于排序 return news; } }提示publish_time是毫秒級(jí)時(shí)間戳不是new Date()字符串。Spark Streaming 的window和slideDuration必須基于事件時(shí)間Event Time否則窗口切分全亂——這是答辯高頻扣分點(diǎn)。源碼里NewsSimulator生成的每條 JSON 都帶這個(gè)字段后續(xù)KafkaProducer發(fā)送時(shí)作為value的一部分Spark 消費(fèi)后可直接.select($publish_time.cast(timestamp))轉(zhuǎn)成時(shí)間列。2.2 Kafka 生產(chǎn)端配置acksallretries3保消息不丟但必須配linger.ms5防小包堆積模擬器生成的 JSON 由KafkaProducer推送到news-topic主題。關(guān)鍵不在代碼多炫酷而在參數(shù)是否符合“畢設(shè)級(jí)可靠性”// 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 副本寫(xiě)成功才返回防 Leader 掛掉丟數(shù)據(jù) props.put(retries, 3); // 網(wǎng)絡(luò)抖動(dòng)時(shí)重試避免單次失敗中斷流 props.put(linger.ms, 5); // 等 5ms 攢一批再發(fā)平衡延遲與吞吐畢設(shè)演示用 5ms 剛好 props.put(batch.size, 16384); // 16KB 批大小配合 linger.ms 防小包泛濫 props.put(buffer.memory, 33554432); // 32MB 緩沖區(qū)防生產(chǎn)過(guò)快阻塞參數(shù)說(shuō)明acksall是必須項(xiàng)答辯時(shí)老師會(huì)問(wèn)“如果 Kafka Leader 宕機(jī)你的消息會(huì)不會(huì)丟”答“默認(rèn) acks1”直接不及格linger.ms5是玄學(xué)調(diào)參點(diǎn)設(shè)為 0 則每條消息單獨(dú)發(fā)網(wǎng)絡(luò)開(kāi)銷(xiāo)大Spark 消費(fèi)端看到的是大量小批次RDD太碎設(shè)為 100ms 則演示時(shí)明顯卡頓等太久。5ms 是實(shí)測(cè)平衡點(diǎn)既保證演示流暢又讓batch.size生效buffer.memory必須大于batch.size * 10否則send()會(huì)阻塞模擬器線程卡死——這是新手常翻車(chē)的黑匣子。2.3 Kafka Topic 創(chuàng)建腳本分區(qū)數(shù)設(shè)為 3副本因子為 2匹配 Spark Streaming 的并行度源碼包里scripts/create_topic.sh直接創(chuàng)建 topic參數(shù)不是隨便寫(xiě)的#!/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 分區(qū) 2 副本Spark Streaming 的KafkaUtils.createDirectStream每個(gè)分區(qū)對(duì)應(yīng)一個(gè) RDD partition分區(qū)數(shù)并行度。3 分區(qū)意味著消費(fèi)端最多啟動(dòng) 3 個(gè) task 并行處理既不過(guò)載單核 CPU又足夠演示“并行消費(fèi)”副本因子 2 保證 ZooKeeper 掛一個(gè)節(jié)點(diǎn)topic 還能讀寫(xiě)——畢設(shè)演示環(huán)境常因虛擬機(jī)內(nèi)存不足 kill 進(jìn)程副本是后悔藥如果你集群只有 1 臺(tái)機(jī)器畢設(shè)常見(jiàn)replication-factor必須 ≤ broker 數(shù)否則create topic報(bào)錯(cuò)Not enough replicas to assign源碼已預(yù)設(shè)為 2你只需確認(rèn)server.properties里broker.id0且listenersPLAINTEXT://:9092正確。3. Spark Streaming 實(shí)時(shí)處理核心窗口聚合不是reduceByKeyAndWindow一把梭而是mapWithState管理熱點(diǎn)詞生命周期3.1 流處理主類NewsStreamingApp的三層架構(gòu)Input → Process → Output拒絕上帝類源碼中com.example.streaming.NewsStreamingApp是入口但它不做任何業(yè)務(wù)邏輯只負(fù)責(zé)組裝// 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: 寫(xiě) MySQL 更新 Redis 緩存 resultStream.foreachRDD { rdd if (!rdd.isEmpty()) { MySQLDao.saveHotWords(rdd) RedisCache.updateTop10(rdd) } } ssc.start() ssc.awaitTermination() }關(guān)鍵設(shè)計(jì)點(diǎn)batchDuration5s是刻意設(shè)的——比 Kafkalinger.ms5略大確保每個(gè) batch 至少攢夠 1~2 條消息避免空 RDD 頻繁觸發(fā)foreachRDDNewsProcessor.process()是獨(dú)立類把“解析 JSON → 提取標(biāo)題關(guān)鍵詞 → 按城市主題窗口統(tǒng)計(jì)”全包進(jìn)去方便單元測(cè)試和答辯時(shí)講清模塊職責(zé)foreachRDD里加if (!rdd.isEmpty())判空否則空 batch 也會(huì)執(zhí)行saveHotWords導(dǎo)致 MySQL 寫(xiě)入空記錄——這是答辯演示時(shí)數(shù)據(jù)表突然多出 10 條NULL的根源。3.2 熱點(diǎn)詞統(tǒng)計(jì)的兩種窗口策略滑動(dòng)窗口保實(shí)時(shí)性mapWithState保狀態(tài)一致性源碼提供兩套統(tǒng)計(jì)方案NewsProcessor默認(rèn)用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) // 注意源碼實(shí)際從 title 提取此處簡(jiǎn)化示意 (s$city-$topic, 1) } // 方案1mapWithState —— 推薦狀態(tài)保存在 checkpoint 目錄重啟不丟計(jì)數(shù) 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 —— 簡(jiǎn)單但狀態(tài)不持久 // parsed.reduceByKeyAndWindow(_ _, _ - _, Minutes(10), Seconds(30)) }為什么mapWithState是畢設(shè)優(yōu)選reduceByKeyAndWindow的窗口狀態(tài)存在內(nèi)存Spark Streaming 重啟后歸零答辯老師問(wèn)“系統(tǒng)掛了 10 分鐘恢復(fù)后熱點(diǎn)詞排名還準(zhǔn)嗎”你答“不準(zhǔn)”就危險(xiǎn)mapWithState的狀態(tài)序列化到 HDFS 或本地checkpoint目錄源碼配在conf/spark-streaming.conf重啟后自動(dòng)加載timeoutMinutes(10)表示 key 10 分鐘沒(méi)更新就自動(dòng)過(guò)期防內(nèi)存泄漏numPartitions(3)強(qiáng)制設(shè)為 3匹配 Kafka 的 3 分區(qū)避免 shuffle——這是性能關(guān)鍵不設(shè)的話默認(rèn)parallelism200小數(shù)據(jù)量反而慢。3.3 關(guān)鍵詞提取的輕量級(jí)實(shí)現(xiàn)不用 jieba 或 HanLP用正則 詞典雙保險(xiǎn)新聞標(biāo)題關(guān)鍵詞提取沒(méi)上 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(發(fā)布, 引發(fā), 熱議, 新政)); private static final String TOPIC_REGEX (AI|新能源|芯片|教育|醫(yī)療); // 預(yù)定義主題詞典 public static ListString extract(String title) { ListString keywords new ArrayList(); // Step1正則匹配預(yù)定義主題詞 Matcher m Pattern.compile(TOPIC_REGEX).matcher(title); while (m.find()) { keywords.add(m.group()); } // Step2按頓號(hào)、逗號(hào)分割過(guò)濾停用詞 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()); } }設(shè)計(jì)理由畢設(shè)不考 NLP 深度考的是“能否在資源受限下合理選型”。jieba 需 Python 環(huán)境HanLP 依賴大詞典而正則詞典 100 行代碼搞定mvn package一鍵打包limit(3)控制每條新聞最多提 3 個(gè)詞防flatMap后數(shù)據(jù)爆炸——Spark Streaming 最怕map后flatMap出 100 倍數(shù)據(jù)executor memory不夠直接 OOM停用詞列表STOP_WORDS可在conf/stopwords.txt動(dòng)態(tài)維護(hù)答辯時(shí)老師問(wèn)“怎么擴(kuò)展新詞”你打開(kāi)文件現(xiàn)場(chǎng)加一行就行。4. 結(jié)果持久化與 Web 展示MySQL 寫(xiě)入不是foreach插入而是JDBCWriter批量 upsert4.1 MySQL DAO 層用PreparedStatement批量 upsert避免 5000 條 insert 變 5000 次網(wǎng)絡(luò)往返com.example.dao.MySQLDao不用 Hibernate 或 MyBatis純 JDBC 手寫(xiě)因?yàn)楫呍O(shè)要體現(xiàn)“底層可控”// 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); // 批量執(zhí)行 for (Tuple2String, Integer t : partition) { ps.setString(1, t._1()); ps.setInt(2, t._2()); ps.setLong(3, System.currentTimeMillis()); ps.addBatch(); // 關(guān)鍵攢批 } ps.executeBatch(); // 一次提交 } finally { if (ps ! null) ps.close(); if (conn ! null) conn.close(); } }); } }參數(shù)與避坑點(diǎn)ON DUPLICATE KEY UPDATE是核心hot_words表主鍵是city_topic重復(fù)插入自動(dòng)累加count避免SELECT INSERT/UPDATE的并發(fā)沖突ps.addBatch()ps.executeBatch()是性能命脈。如果寫(xiě)成ps.executeUpdate()循環(huán)1000 條數(shù)據(jù)就是 1000 次 TCP 握手演示時(shí)卡成 PPTuseSSLfalse必須顯式加MySQL 5.7 默認(rèn) require SSL不加這參數(shù)getConnection直接拋SQLException——這是源碼包能跑通但你本地跑不通的頭號(hào)原因。4.2 Spring Boot Web 層REST 接口返回 JSON不渲染頁(yè)面用curl或 Postman 即可驗(yàn)證com.example.web.HotWordController只提供兩個(gè)接口極簡(jiǎn)// 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); } }部署要點(diǎn)application.properties里server.port8081避開(kāi)了 Spark History Server 的 18080 和 Kafka 的 9092HotWord實(shí)體類用DataLombokmvn clean package前確保 IDEA 安裝 Lombok 插件否則編譯報(bào)錯(cuò)——這是新手 IDE 環(huán)境沒(méi)配好的典型翻車(chē)接口返回ListHotWord前端用fetch(/api/top10).then(r r.json())即可源碼包web/static/index.html里已寫(xiě)好輪詢 JSsetInterval(() loadTop10(), 3000)3 秒刷新一次演示時(shí)效果直觀。4.3 數(shù)據(jù)庫(kù)建表 SQLcity_topic設(shè)為主鍵count加索引last_update用TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMPscripts/init_mysql.sql腳本建表字段設(shè)計(jì)直指痛點(diǎn)-- 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 排序時(shí)走索引 ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;為什么這樣建city_topic作主鍵天然去重INSERT ... ON DUPLICATE KEY UPDATE才生效INDEX idx_count (count)是給ORDER BY count DESC LIMIT 10用的。沒(méi)這個(gè)索引queryTop10()會(huì)全表掃描10 萬(wàn)行數(shù)據(jù)查 10 條要 2 秒演示時(shí)卡頓utf8mb4支持 emoji 和生僻字新聞標(biāo)題可能含“ SpaceX”或“量子”等字符utf8會(huì)亂碼——答辯老師看到亂碼直接質(zhì)疑數(shù)據(jù)質(zhì)量。5. 避坑指南Spark 2.2 Kafka 0.10.2 組合下90% 的失敗都源于這 5 個(gè)具體錯(cuò)誤5.1 現(xiàn)象ClassNotFoundException: org.apache.spark.streaming.kafka010.KafkaUtils原因spark-streaming-kafka-0-10_2.11依賴 scope 寫(xiě)成了compile導(dǎo)致打包時(shí)把 Kafka client 打進(jìn) 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表示運(yùn)行時(shí)由集群提供mvn package不打進(jìn) jar。本地測(cè)試用mvn compile exec:java集群提交用spark-submit --jars /path/to/spark-streaming-kafka-0-10_2.11-2.2.0.jar顯式指定。5.2 現(xiàn)象Spark Streaming 啟動(dòng)后無(wú)日志ssc.awaitTermination()一直阻塞原因spark.streaming.kafka.consumer.poll.ms默認(rèn) 512ms但 Kafka broker 地址配錯(cuò)如localhost:9092在集群模式下應(yīng)為kafka-host:9092consumer 無(wú)法連接poll()永久超時(shí)。解決在conf/spark-streaming.conf中顯式設(shè)置超時(shí)并打印 debug 日志spark.streaming.kafka.consumer.poll.ms1000 spark.logConftrue spark.eventLog.enabledtrue然后spark-submit加--conf spark.driver.extraJavaOptions-Dlog4j.debug查看 Kafka 連接日志。5.3 現(xiàn)象MySQL 表里count字段突增 100 倍數(shù)據(jù)嚴(yán)重失真原因NewsSimulator生成速度過(guò)快如ScheduledExecutorService周期設(shè)為 100ms而 Spark batch interval 為 5s導(dǎo)致單個(gè) batch 拉取到 50 條消息mapWithState對(duì)同一city-topickey 累加 50 次但業(yè)務(wù)邏輯本意是“每條新聞貢獻(xiàn) 1 熱度”。解決嚴(yán)格控制模擬器吞吐量。NewsSimulator.java中scheduleAtFixedRate的 delay 設(shè)為20002 秒確保 5s batch 內(nèi)最多 2~3 條新聞。演示前用kafka-console-consumer.sh手動(dòng)驗(yàn)證消息速率。5.4 現(xiàn)象Web 頁(yè)面top10接口返回空數(shù)組但 MySQL 表里有數(shù)據(jù)原因MySQLDao.queryTop10()方法里PreparedStatement的ORDER BY count DESC LIMIT 10沒(méi)加rs.next()判空ResultSet 為空時(shí)while(rs.next())不執(zhí)行l(wèi)ist 返回空。解決源碼中已修復(fù)但若你修改過(guò) DAO請(qǐng)確保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 現(xiàn)象spark-submit報(bào)java.lang.NoClassDefFoundError: scala/Product原因Scala 版本不匹配。Spark 2.2 編譯于 Scala 2.11但你的pom.xml里scala.version設(shè)為 2.12。解決統(tǒng)一 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. 畢設(shè)答辯加分技巧用spark-shell實(shí)時(shí)驗(yàn)證流處理邏輯三步定位 Kafka 消費(fèi)斷點(diǎn)6.1 第一步用spark-shell直連 Kafka驗(yàn)證消息是否真的在流里別等NewsStreamingApp啟動(dòng)失敗才排查先用 Spark Shell 快速驗(yàn)貨$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]進(jìn)入 shell 后執(zhí)行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秒后自動(dòng)停作用如果這 5 行能打印出 JSON證明 Kafka 消息正常如果卡住或報(bào)錯(cuò)問(wèn)題在 Kafka 配置不是 Spark 代碼。這招能在答辯前 1 小時(shí)快速鎖定故障域。6.2 第二步用kafka-console-consumer查看 offset確認(rèn) consumer group 是否滯留Spark Streaming 的group.id默認(rèn)是隨機(jī)的但調(diào)試時(shí)需固定# 查看 news-topic 的所有 consumer group $KAFKA_HOME/bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --list # 查看指定 group 的 offset假設(shè) group.iddebug-group $KAFKA_HOME/bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \ --group debug-group --describe --topic news-topic關(guān)鍵指標(biāo)CURRENT-OFFSET當(dāng)前消費(fèi)到的位置LOG-END-OFFSET最新消息位置LAG差值即積壓消息數(shù)。如果LAG 0且持續(xù)增長(zhǎng)說(shuō)明 Spark 消費(fèi)太慢需調(diào)大 executor cores 或減小 batchDuration。6.3 第三步在NewsProcessor中加printlncheckpoint目錄監(jiān)控確認(rèn)狀態(tài)是否真正持久化mapWithState的狀態(tài)存在checkpoint目錄但默認(rèn)路徑是/tmp/checkpoint重啟后可能被清空。源碼中已配為hdfs://localhost:9000/checkpoint或本地file:///opt/spark-checkpoint你必須確認(rèn)conf/spark-streaming.conf中spark.streaming.checkpointDirectory路徑可寫(xiě)啟動(dòng)NewsStreamingApp后立刻ls -l /opt/spark-checkpoint應(yīng)看到00000000000000000000等數(shù)字目錄手動(dòng)kill -9進(jìn)程再啟動(dòng)觀察日志是否有Loading state from checkpoint字樣。表格checkpoint 目錄關(guān)鍵文件含義文件名作用答辯時(shí)可展示offsets/存 Kafka 每個(gè) partition 的消費(fèi) offset保證 Exactly-Oncecat offsets/00000000000000000000顯示 offset 值state/mapWithState的狀態(tài)快照二進(jìn)制格式du -sh state/證明狀態(tài)非空streaming/StreamingContext 元數(shù)據(jù)ls streaming/證明 checkpoint 目錄被識(shí)別從那以后我每次重構(gòu)NewsProcessor都強(qiáng)制走一遍spark-shell驗(yàn)證流、kafka-consumer-groups查 offset、ls checkpoint看狀態(tài)三步曲——不是為了炫技是避免答辯前夜發(fā)現(xiàn)mapWithState根本沒(méi)生效只能重寫(xiě)邏輯。這套源碼的真正價(jià)值不在于它多完美而在于它把 Spark Streaming 從黑匣子拆成了可觸摸、可驗(yàn)證、可答辯的零件。希望幫到你。本文還有配套的精品資源點(diǎn)擊獲取