Spark Streaming微批次架構(gòu)解析與實時計算實踐指南
1. 項目概述為什么Spark Streaming依然是實時計算的基石最近和幾個做數(shù)據(jù)平臺的朋友聊天發(fā)現(xiàn)一個挺有意思的現(xiàn)象盡管現(xiàn)在實時計算領(lǐng)域新框架層出不窮比如Flink風頭正勁但在很多公司的生產(chǎn)環(huán)境里Spark Streaming依然穩(wěn)穩(wěn)地占據(jù)著一席之地處理著大量的實時數(shù)據(jù)流。這讓我想起了自己幾年前第一次接觸Spark Streaming的場景當時為了搞定一個簡單的實時點擊流統(tǒng)計折騰了好幾個晚上。現(xiàn)在回過頭看Spark Streaming的設(shè)計理念其實非常經(jīng)典它把流處理巧妙地“偽裝”成了一系列連續(xù)的微批次Micro-Batch處理這種“流批一體”的早期思想讓很多熟悉Spark批處理Spark Core的開發(fā)者能夠幾乎無門檻地上手實時計算。簡單來說Spark Streaming是Apache Spark生態(tài)系統(tǒng)里用于處理實時數(shù)據(jù)流的組件。它的核心能力是能夠從Kafka、Flume、Kinesis或者TCP Socket等多種數(shù)據(jù)源接入高速數(shù)據(jù)流然后利用Spark強大的分布式計算引擎對這些數(shù)據(jù)進行高吞吐、可容錯的實時處理最后將結(jié)果輸出到文件系統(tǒng)、數(shù)據(jù)庫或者實時儀表盤。它解決的痛點很明確在數(shù)據(jù)產(chǎn)生的瞬間就進行分析和響應而不是等到攢夠一批再處理這對于監(jiān)控、風控、實時推薦等場景至關(guān)重要。那么誰適合深入了解一下Spark Streaming呢如果你已經(jīng)是Spark批處理的用戶想將業(yè)務擴展到實時領(lǐng)域那么Spark Streaming是你的自然選擇學習曲線非常平緩。如果你在評估實時計算框架需要的是一個成熟、穩(wěn)定、社區(qū)資源豐富并且能與現(xiàn)有Spark批處理作業(yè)無縫整合的方案Spark Streaming也值得你重點考察。當然對于初學者而言理解Spark Streaming的微批次模型也是理解現(xiàn)代流處理編程范式的絕佳起點。接下來我會結(jié)合自己踩過的坑和積累的經(jīng)驗帶你從設(shè)計思路到實操細節(jié)徹底搞懂Spark Streaming。2. 核心架構(gòu)與微批次模型深度解析2.1 DStream流計算的核心抽象Spark Streaming的編程模型核心是離散化流也就是DStream。這是理解其一切行為的關(guān)鍵。很多新手會困惑為什么我的流處理作業(yè)延遲感覺不像Flink那么“實時”答案就藏在DStream的設(shè)計里。你可以把DStream想象成一個連續(xù)不斷的“數(shù)據(jù)序列”但這個序列不是平滑的而是被切成了一個個固定時間間隔的“數(shù)據(jù)切片”。每一個切片本質(zhì)上就是一個RDD彈性分布式數(shù)據(jù)集。也就是說一個DStream在背后是由一系列按時間順序排列的RDD所構(gòu)成的。Spark Streaming的作業(yè)調(diào)度器會周期性地這個周期就是你設(shè)置的批次間隔比如1秒啟動Spark作業(yè)來處理當前時間窗口內(nèi)到達的、屬于同一個RDD的數(shù)據(jù)。舉個例子你設(shè)置批次間隔為2秒。那么Spark Streaming會每2秒創(chuàng)建一個新的RDD這個RDD包含了這2秒內(nèi)從數(shù)據(jù)源接收到的所有數(shù)據(jù)。然后你定義的所有轉(zhuǎn)換操作如map、filter、reduceByKey都會作用在這個RDD上生成新的DStream。這種設(shè)計帶來了幾個深遠的影響與Spark Core的無縫繼承所有你在批處理中熟悉的RDD操作、持久化、容錯機制在DStream上幾乎完全適用。你的知識復用率極高。一致的語義因為底層是RDD所以Spark Streaming能提供“精確一次”的語義保障這對于金融、交易類場景是硬性要求。這通常需要與可靠的數(shù)據(jù)源如Kafka Direct API和可靠的輸出協(xié)同工作。吞吐量優(yōu)先微批次模型天生有利于吞吐量。它可以將一小段時間內(nèi)的數(shù)據(jù)攢起來進行優(yōu)化后再計算非常適合高吞吐的日志處理、指標聚合場景。注意這個“批次間隔”是你調(diào)優(yōu)的第一個關(guān)鍵參數(shù)。設(shè)置得太短如100ms會導致調(diào)度開銷過大可能每個批次的數(shù)據(jù)量很小無法充分發(fā)揮集群性能設(shè)置得太長如10秒又會導致數(shù)據(jù)處理延遲變高實時性變差。通常在生產(chǎn)環(huán)境中1-5秒是一個常見的起始探索區(qū)間。2.2 容錯與狀態(tài)管理機制流處理系統(tǒng)必須可靠。Spark Streaming的容錯建立在RDD的血統(tǒng)Lineage機制之上。每個RDD都知道它是如何從父RDD計算而來的。如果某個節(jié)點宕機導致某個RDD分區(qū)丟失Spark可以直接根據(jù)血統(tǒng)重新計算該分區(qū)從而實現(xiàn)數(shù)據(jù)恢復。但對于有狀態(tài)的計算例如計算最近10分鐘的用戶點擊次數(shù)僅僅重新計算丟失的數(shù)據(jù)是不夠的因為狀態(tài)本身可能已經(jīng)累積了很久。為此Spark Streaming引入了檢查點機制和狀態(tài)DStream。檢查點有兩種類型。元數(shù)據(jù)檢查點將流計算應用的DAG信息、配置等持久化到HDFS等可靠存儲用于驅(qū)動程序的故障恢復。如果你的Driver程序掛掉重啟后可以從檢查點恢復上下文并繼續(xù)處理。數(shù)據(jù)檢查點將中間生成的RDD定期保存。這對于那些血統(tǒng)鏈過長例如使用了updateStateByKey且窗口很大的DStream尤為重要可以切斷過長的依賴鏈避免恢復時重新計算整個歷史。狀態(tài)管理對于需要跨批次維護狀態(tài)的操作早期主要使用updateStateByKey。它允許你為每個Key維護一個任意類型的狀態(tài)并在每個批次更新它。但這個方法有個問題它會在每個批次都對所有Key進行計算即使這個Key在本批次沒有新數(shù)據(jù)這在小批次間隔下會帶來不小的開銷。后來Spark引入了更高效的mapWithStateAPI。它只對那些在本批次有更新的Key進行狀態(tài)更新和輸出性能提升非常顯著。在最新的Structured Streaming中狀態(tài)管理得到了進一步的抽象和優(yōu)化。2.3 與Structured Streaming的關(guān)系辨析這是當前Spark流處理生態(tài)中一個必須厘清的概念。Structured Streaming是Spark 2.0后引入的新的流處理引擎它不再基于DStream而是基于Spark SQL引擎將數(shù)據(jù)流視為一張無限增長的表。特性Spark Streaming (DStreams)Structured Streaming編程模型基于RDD的底層API基于DataFrame/Dataset的高級APIAPI級別相對底層靈活性高聲明式更高級更簡潔時間語義主要處理處理時間原生支持事件時間、處理時間以及延遲數(shù)據(jù)的處理水位線不支持支持用于處理亂序事件狀態(tài)管理updateStateByKey/mapWithState內(nèi)建支持更簡單容錯語義可達到精確一次端到端精確一次需配合特定Source/Sink與批處理統(tǒng)一共享RDD API共享DataFrame API真正做到代碼統(tǒng)一如何選擇對于新項目強烈建議優(yōu)先考慮Structured Streaming。它在易用性、時間語義支持和與批處理的統(tǒng)一性上優(yōu)勢明顯。那為什么還要學DStream呢首先大量遺留系統(tǒng)仍在運行Spark Streaming維護和優(yōu)化需要相關(guān)知識。其次DStream API讓你更接近底層對于理解流計算的本質(zhì)、進行一些極其定制化的操作雖然很少需要仍有價值。最后學習DStream的微批次模型能幫你更好地理解Structured Streaming在底層是如何工作的。3. 從零到一一個完整的Spark Streaming應用實戰(zhàn)理論說得再多不如動手跑一遍。我們來實現(xiàn)一個經(jīng)典的場景從Kafka讀取用戶行為日志JSON格式實時統(tǒng)計每10秒內(nèi)每個頁面的訪問量PV并將結(jié)果輸出到控制臺和MySQL數(shù)據(jù)庫。3.1 環(huán)境準備與依賴配置首先你需要一個Spark環(huán)境。本地測試最簡單的方式是下載Spark預編譯包解壓即可。生產(chǎn)環(huán)境則通常部署在YARN或Kubernetes上。我們假設(shè)使用本地模式進行演示。創(chuàng)建一個標準的Maven或SBT項目。關(guān)鍵的依賴包括!-- Spark Streaming 核心 -- dependency groupIdorg.apache.spark/groupId artifactIdspark-streaming_2.12/artifactId version3.3.0/version !-- 請使用與Spark Core一致的版本 -- /dependency !-- 用于連接Kafka -- dependency groupIdorg.apache.spark/groupId artifactIdspark-streaming-kafka-0-10_2.12/artifactId version3.3.0/version /dependency !-- MySQL連接器用于輸出 -- dependency groupIdmysql/groupId artifactIdmysql-connector-java/artifactId version8.0.33/version /dependency實操心得依賴的Scala版本這里是2.12必須與你安裝的Spark運行時版本嚴格一致否則會引發(fā)各種詭異的NoSuchMethodError。最好通過spark-shell --version命令確認你的Spark環(huán)境版本。3.2 應用主邏輯編寫下面是完整的Scala應用示例。我們使用Kafka的Direct API無Receiver模式這是目前推薦的方式具有更好的并行度和一致性語義。import org.apache.spark.SparkConf import org.apache.spark.streaming.{Seconds, StreamingContext} import org.apache.spark.streaming.kafka010._ import org.apache.spark.streaming.kafka010.LocationStrategies.PreferConsistent import org.apache.spark.streaming.kafka010.ConsumerStrategies.Subscribe import org.json4s._ import org.json4s.jackson.JsonMethods._ import java.sql.{Connection, DriverManager, PreparedStatement} import java.util.Properties object RealtimePageViewCounter { // 隱式參數(shù)用于json4s解析 implicit val formats: DefaultFormats DefaultFormats case class UserLog(userId: String, pageId: String, timestamp: Long) def main(args: Array[String]): Unit { // 1. 創(chuàng)建SparkConf和StreamingContext批次間隔設(shè)為2秒 val sparkConf new SparkConf() .setAppName(RealtimePageViewCounter) .setMaster(local[2]) // 本地測試用2個核生產(chǎn)環(huán)境去掉此參數(shù)通過spark-submit指定 .set(spark.serializer, org.apache.spark.serializer.KryoSerializer) // 使用Kryo序列化提升性能 val ssc new StreamingContext(sparkConf, Seconds(2)) // 設(shè)置檢查點目錄用于狀態(tài)恢復本地測試可先注釋 // ssc.checkpoint(hdfs://your-nn:9000/spark-streaming-checkpoint) // 2. 配置Kafka參數(shù) val kafkaParams Map[String, Object]( bootstrap.servers - kafka-broker1:9092,kafka-broker2:9092, key.deserializer - org.apache.kafka.common.serialization.StringDeserializer, value.deserializer - org.apache.kafka.common.serialization.StringDeserializer, group.id - spark-streaming-pageview-group, auto.offset.reset - latest, // 從最新位置開始消費 enable.auto.commit - (false: java.lang.Boolean) // Spark自己管理offset ) val topics Array(user-behavior-topic) // 3. 創(chuàng)建DStream連接Kafka val stream KafkaUtils.createDirectStream[String, String]( ssc, PreferConsistent, Subscribe[String, String](topics, kafkaParams) ) // 4. 數(shù)據(jù)處理邏輯 val pageCounts stream .map(record record.value()) // 提取Kafka消息的值JSON字符串 .filter(_.nonEmpty) // 過濾空消息 .map { jsonString try { // 解析JSON提取pageId val json parse(jsonString) val pageId (json \ pageId).extractOrElse[String](unknown) (pageId, 1) } catch { case e: Exception // 記錄解析錯誤實際生產(chǎn)中應寫入錯誤日志或死信隊列 println(sFailed to parse JSON: $jsonString, error: ${e.getMessage}) (parse_error, 1) } } .reduceByKeyAndWindow( _ _, // 聚合函數(shù)累加 _ - _, // 逆函數(shù)用于窗口滑動時減去過期批次提升性能需設(shè)置檢查點 Seconds(10), // 窗口長度10秒 Seconds(2) // 滑動間隔2秒與批次間隔相同 ) // 如果不使用逆函數(shù)可以用簡單的 reduceByKey(_ _).window(Seconds(10), Seconds(2)) // 5. 輸出操作觸發(fā)計算并輸出 pageCounts.foreachRDD { (rdd, time) // 注意foreachRDD內(nèi)部的代碼在Driver端執(zhí)行但其中的RDD操作在Executor端執(zhí)行 if (!rdd.isEmpty()) { println(s\n Batch Time: $time ) // 輸出到控制臺 rdd.foreachPartition { partitionOfRecords // 這個foreach在Executor上執(zhí)行 partitionOfRecords.foreach { case (pageId, count) println(sPage: $pageId, Count: $count) } } // 輸出到MySQL (在Driver端收集少量數(shù)據(jù)后寫入或使用foreachPartition在Executor寫) // 方式A收集到Driver后寫入適合結(jié)果集小 val collectedData rdd.collect() if (collectedData.nonEmpty) { saveToMySQL(collectedData, time) } // 方式B使用foreachPartition在Executor分布式寫入適合結(jié)果集大但需管理連接池 // rdd.foreachPartition { partition // val conn getMySQLConnection() // // ... 批量插入邏輯 // conn.close() // } } } // 6. 啟動流計算并等待終止 ssc.start() ssc.awaitTermination() } def saveToMySQL(data: Array[(String, Int)], batchTime: org.apache.spark.streaming.Time): Unit { var conn: Connection null var pstmt: PreparedStatement null val url jdbc:mysql://your-mysql-host:3306/streaming_db val user your_user val password your_password try { Class.forName(com.mysql.cj.jdbc.Driver) conn DriverManager.getConnection(url, user, password) // 假設(shè)表結(jié)構(gòu)page_pv (batch_time TIMESTAMP, page_id VARCHAR(50), pv INT) val sql INSERT INTO page_pv (batch_time, page_id, pv) VALUES (?, ?, ?) ON DUPLICATE KEY UPDATE pv ? pstmt conn.prepareStatement(sql) val timestamp new java.sql.Timestamp(batchTime.milliseconds) for ((pageId, count) - data) { pstmt.setTimestamp(1, timestamp) pstmt.setString(2, pageId) pstmt.setInt(3, count) pstmt.setInt(4, count) // 用于ON DUPLICATE KEY UPDATE pstmt.addBatch() } pstmt.executeBatch() println(sSuccessfully saved ${data.length} records to MySQL.) } catch { case e: Exception e.printStackTrace() } finally { if (pstmt ! null) pstmt.close() if (conn ! null) conn.close() } } }3.3 關(guān)鍵代碼段解析與調(diào)優(yōu)點StreamingContext初始化這是所有流計算的起點。Seconds(2)定義了微批次的間隔。local[2]中的數(shù)字2代表至少使用2個CPU核心一個用于接收數(shù)據(jù)一個用于處理數(shù)據(jù)這是本地測試的最低要求。Kafka Direct API我們使用createDirectStream并設(shè)置enable.auto.commit為false。這意味著Spark Streaming會自己將消費偏移量offset管理在檢查點中或自己提交回Kafka這是實現(xiàn)“精確一次”處理的基礎(chǔ)。你需要確保輸出操作是冪等的或者將offset和輸出結(jié)果放在同一個事務中。reduceByKeyAndWindow這是窗口操作的核心。我們設(shè)置了10秒的窗口長度和2秒的滑動間隔。這意味著每2秒一個批次我們會計算過去10秒內(nèi)的數(shù)據(jù)。使用了加法和減法函數(shù)這要求開啟檢查點但能極大優(yōu)化滑動窗口的性能因為它不需要重復計算重疊部分的數(shù)據(jù)。foreachRDD的設(shè)計模式這是輸出結(jié)果到外部系統(tǒng)如數(shù)據(jù)庫、Redis的標準入口。至關(guān)重要的一點foreachRDD內(nèi)部的代碼在Driver端執(zhí)行但其中的RDD操作如foreachPartition是在Executor端執(zhí)行的。創(chuàng)建數(shù)據(jù)庫連接等昂貴操作應該在foreachPartition內(nèi)部進行并為每個分區(qū)創(chuàng)建一個連接池而不是為每條記錄創(chuàng)建連接更不要在Driver端創(chuàng)建連接然后序列化到Executor這會導致序列化錯誤。4. 生產(chǎn)環(huán)境部署與性能調(diào)優(yōu)指南把應用跑起來只是第一步要讓它在生產(chǎn)環(huán)境中穩(wěn)定、高效地運行還需要做大量工作。4.1 資源分配與并行度優(yōu)化Spark Streaming應用的性能很大程度上取決于資源是否給夠以及任務是否被充分并行化。Executor資源通過spark-submit提交時需要合理設(shè)置。--num-executorsExecutor數(shù)量。根據(jù)數(shù)據(jù)量和處理邏輯復雜度決定通常從10-20個開始。--executor-cores每個Executor的CPU核心數(shù)。建議2-4個確保每個Executor能并行執(zhí)行多個任務。--executor-memory每個Executor的內(nèi)存。需要容納接收到的批次數(shù)據(jù)、進行轉(zhuǎn)換操作產(chǎn)生的中間數(shù)據(jù)以及維護的狀態(tài)。必須預留一部分給操作系統(tǒng)和HDFS客戶端約10%。例如總內(nèi)存4G可設(shè)置--executor-memory 3g。Receiver與并行度如果使用舊的Receiver模式不推薦接收數(shù)據(jù)本身會占用一個CPU核心。在Direct API下Kafka分區(qū)數(shù)直接決定了讀取階段的并行度。確保Kafka主題的分區(qū)數(shù) Spark Streaming作業(yè)中讀取該主題的并發(fā)任務數(shù)。通常你可以通過spark.streaming.kafka.maxRatePerPartition參數(shù)控制每個分區(qū)每秒讀取的最大消息數(shù)來平衡吞吐和延遲。處理并行度由RDD的分區(qū)數(shù)決定。Shuffle操作如reduceByKey后的默認分區(qū)數(shù)由spark.default.parallelism控制通常設(shè)置為executor-cores * num-executors的2-3倍。你也可以在操作中顯式指定分區(qū)數(shù)如reduceByKey(__, 100)。4.2 背壓機制與動態(tài)資源分配當數(shù)據(jù)流入速度超過處理速度時會導致批次處理時間越來越長最終堆積崩潰。Spark Streaming 1.5之后引入了背壓機制可以動態(tài)調(diào)整接收速率來適配處理能力。啟用背壓設(shè)置spark.streaming.backpressure.enabledtrue。背壓算法默認使用PID控制器你也可以通過spark.streaming.backpressure.initialRate設(shè)置初始接收速率。啟用后Spark會監(jiān)控批次處理時間和調(diào)度延遲自動調(diào)整從Kafka等源拉取數(shù)據(jù)的速率。對于運行在YARN上的應用還可以結(jié)合動態(tài)資源分配。但這在流處理中需謹慎使用因為申請和釋放Executor需要時間可能影響實時性。通常適用于處理負載有明顯波峰波谷且對延遲不極度敏感的場景。4.3 檢查點與狀態(tài)恢復實戰(zhàn)檢查點不是可選項對于生產(chǎn)應用是必選項。它用于元數(shù)據(jù)恢復和狀態(tài)計算。設(shè)置檢查點目錄目錄必須是一個可靠的文件系統(tǒng)如HDFS。ssc.checkpoint(“hdfs://...”。編寫可恢復的驅(qū)動程序你的主函數(shù)需要能被Spark在故障后重新調(diào)用。標準模式如下def createStreamingContext(): StreamingContext { val sparkConf ... val ssc new StreamingContext(sparkConf, Seconds(2)) // 定義你的DStream計算邏輯 val lines ... // ... ssc.checkpoint(checkpointDir) ssc } def main(args: Array[String]) { val checkpointDir “hdfs://...” val ssc StreamingContext.getOrCreate(checkpointDir, createStreamingContext _) ssc.start() ssc.awaitTermination() }這樣當Driver重啟時getOrCreate會嘗試從檢查點目錄重建StreamingContext。如果失敗則調(diào)用提供的函數(shù)創(chuàng)建新的。踩坑實錄檢查點目錄包含了序列化的類。如果你修改了應用代碼如添加了新的類字段然后試圖從舊的檢查點恢復會引發(fā)序列化錯誤。最佳實踐是每次代碼升級后清空檢查點目錄意味著從最新的Kafka偏移量開始消費或者確保代碼變更向后兼容。5. 典型問題排查與監(jiān)控運維即使應用部署成功運維過程中也會遇到各種問題。這里記錄幾個最常見的問題和排查思路。5.1 批次處理延遲與堆積這是最常見的問題。癥狀是Spark UI的Streaming頁面上批次處理時間Processing Time持續(xù)大于批次間隔Batch Interval導致“Scheduling Delay”不斷增長。排查步驟看日志首先查看Executor和Driver的日志是否有明顯的錯誤或GC警告??碨park UIStreaming頁確認哪些批次延遲了。是持續(xù)延遲還是偶發(fā)Stages頁點擊延遲批次對應的作業(yè)查看是哪個Stage耗時最長。是讀取數(shù)據(jù)慢Shuffle慢還是輸出慢Executors頁觀察GC時間是否過長。如果Full GC頻繁說明內(nèi)存不足。針對性優(yōu)化數(shù)據(jù)傾斜如果某個Stage的某個Task執(zhí)行時間遠長于其他很可能是數(shù)據(jù)傾斜。使用sample方法查看Key分布考慮使用加鹽隨機前綴打散熱點Key。外部系統(tǒng)瓶頸如果延遲發(fā)生在foreachRDD的輸出階段可能是數(shù)據(jù)庫或Redis寫入慢??紤]使用連接池、批量寫入、異步寫入或換用更高性能的輸出端。資源不足如果所有Task都慢且GC正常可能是CPU或內(nèi)存整體不足。嘗試增加executor-cores或executor-memory。調(diào)整批次間隔適當增大批次間隔如從1秒到2秒給每個批次更多處理時間可以緩解短期壓力但會犧牲實時性。5.2 數(shù)據(jù)丟失與重復消費這通常與偏移量管理和輸出操作的原子性有關(guān)。確保精確一次語義使用Direct API它讓Spark自己管理Kafka偏移量。可靠的數(shù)據(jù)源確保Kafka本身是高可用的。冪等的輸出或事務性輸出這是最難的部分。要么你的輸出操作是冪等的比如INSERT ON DUPLICATE KEY UPDATE要么你將偏移量的提交和數(shù)據(jù)的輸出放在同一個數(shù)據(jù)庫事務中。Spark本身不提供跨系統(tǒng)的事務這需要你在foreachRDD中自己實現(xiàn)。監(jiān)控偏移量定期檢查Spark提交到Kafka的消費者組偏移量確保其正常推進并與實際處理進度匹配。5.3 監(jiān)控與告警體系搭建不能等用戶投訴了才發(fā)現(xiàn)流處理作業(yè)掛了。必須建立監(jiān)控。Spark UI History Server這是最基本的。通過History Server可以查看已結(jié)束應用的運行情況。Metrics系統(tǒng)Spark提供了豐富的Metrics可以通過SparkConf配置輸出到Ganglia、Graphite、Prometheus等系統(tǒng)。關(guān)鍵指標包括spark.streaming.*: 如processingDelay處理延遲、schedulingDelay調(diào)度延遲、numReceivers接收器數(shù)量、numTotalCompletedBatches總完成批次數(shù)等。JVM相關(guān)指標GC時間、堆內(nèi)存使用情況。自定義應用指標你可以在foreachRDD里將每批次處理的數(shù)據(jù)量、輸出記錄數(shù)等業(yè)務指標推送到你的監(jiān)控系統(tǒng)如StatsD。進程存活監(jiān)控使用系統(tǒng)級的監(jiān)控工具如Supervisord、K8s Liveness Probe確保Driver和Executor進程存活。對于YARN可以監(jiān)控YARN Application狀態(tài)。告警規(guī)則針對關(guān)鍵指標設(shè)置告警例如連續(xù)N個批次處理延遲超過閾值、消費者組滯后Lag持續(xù)增長、Executor頻繁丟失等。Spark Streaming是一個經(jīng)歷過大規(guī)模生產(chǎn)環(huán)境考驗的框架它的微批次模型在吞吐量和一致性之間取得了很好的平衡。雖然Structured Streaming代表了未來的方向但理解DStream的運作機制、掌握其調(diào)優(yōu)和運維技巧對于任何一個大數(shù)據(jù)開發(fā)者來說仍然是一筆寶貴的財富。在實際項目中最關(guān)鍵的是根據(jù)業(yè)務對延遲和吞吐量的具體要求以及對一致性的容忍度來做出最合適的架構(gòu)選擇和技術(shù)決策。

相關(guān)新聞

2026年P(guān)DF轉(zhuǎn)換器實測盤點:免費好用、離線安全、手機端方案一次說清

2026年P(guān)DF轉(zhuǎn)換器實測盤點:免費好用、離線安全、手機端方案一次說清

2026年P(guān)DF轉(zhuǎn)換器實測盤點:免費好用、離線安全、手機端方案一次說清 前陣子同事扔過來一份簽完字的合同掃描件,讓我把關(guān)鍵條款摘出來補進報告。文件不大,但排版密密麻麻,頁腳還有手寫備注。我下意識先掏出手機,在微信里…

2026/8/2 12:56:09 閱讀更多
高質(zhì)量數(shù)據(jù)集的特征

高質(zhì)量數(shù)據(jù)集的特征

通識高質(zhì)量數(shù)據(jù)集、行業(yè)通識高質(zhì)量數(shù)據(jù)集和行業(yè)專識高質(zhì)量數(shù)據(jù)集,都具有知識內(nèi)容、來源類型、時效性、標注人員類型、敏感程度、模型類型、主題范圍等方面的特征。知識內(nèi)容指數(shù)據(jù)集中數(shù)據(jù)所蘊含知識的專業(yè)性、知識深度和目標受眾;來源類型指數(shù)據(jù)集中數(shù)據(jù)的獲取來源&…

2026/8/2 14:06:11 閱讀更多
串口藍牙模塊實戰(zhàn)指南:從原理到Arduino智能小車遙控應用

串口藍牙模塊實戰(zhàn)指南:從原理到Arduino智能小車遙控應用

1. 項目概述:為什么我們還需要一個串口藍牙模塊?在物聯(lián)網(wǎng)和智能硬件項目里,無線通信幾乎是標配。Wi-Fi、藍牙、LoRa、Zigbee……選擇很多。但如果你問一個經(jīng)常搗鼓Arduino、樹莓派或者ESP32的開發(fā)者,哪種無線連接方式最“無腦”、…

2026/8/2 14:06:11 閱讀更多
怎樣用AI魔法讓模糊視頻變清晰:3個簡單秘訣

怎樣用AI魔法讓模糊視頻變清晰:3個簡單秘訣

怎樣用AI魔法讓模糊視頻變清晰:3個簡單秘訣 【免費下載鏈接】video2x A machine learning-based video super resolution and frame interpolation framework. Est. Hack the Valley II, 2018. 項目地址: https://gitcode.com/GitHub_Trending/vi/video2x 你…

2026/8/2 13:56:11 閱讀更多
3分鐘搞定!QQ空間歷史說說完整備份終極指南

3分鐘搞定!QQ空間歷史說說完整備份終極指南

3分鐘搞定!QQ空間歷史說說完整備份終極指南 【免費下載鏈接】GetQzonehistory 獲取QQ空間發(fā)布的歷史說說 項目地址: https://gitcode.com/GitHub_Trending/ge/GetQzonehistory 你是否曾想過,那些年發(fā)過的QQ空間說說,那些記錄青春的文字…

2026/8/2 0:04:01 閱讀更多
3分鐘搞定!QQ空間歷史說說完整備份終極指南

3分鐘搞定!QQ空間歷史說說完整備份終極指南

3分鐘搞定!QQ空間歷史說說完整備份終極指南 【免費下載鏈接】GetQzonehistory 獲取QQ空間發(fā)布的歷史說說 項目地址: https://gitcode.com/GitHub_Trending/ge/GetQzonehistory 你是否曾想過,那些年發(fā)過的QQ空間說說,那些記錄青春的文字…

2026/8/2 0:04:01 閱讀更多
AMAT 0100-02186 I/O 分配 PCB

AMAT 0100-02186 I/O 分配 PCB

AMAT 0100-02186 I/O分配PCB板是應用材料(Applied Materials)公司生產(chǎn)的一款用于半導體設(shè)備的I/O信號分配電路板。該型號(0100-02186)的核心特點如下:專用于Endura等半導體工藝腔室。集成信號路由與分配功能。連接控制…

2026/8/2 2:51:21 閱讀更多
Nissei Corp FFMN-32L-10-T0 40AX 三相異步電動機

Nissei Corp FFMN-32L-10-T0 40AX 三相異步電動機

Nissei Corp FFMN-32L-10-T0 40AX 三相異步電動機是日本日清(Nissei)品牌的一款工業(yè)用三相異步電機,適用于自動化設(shè)備及通用機械驅(qū)動。該型號(FFMN-32L-10-T0 40AX)的核心特點如下:三相交流異步電動機。額定…

2026/8/2 2:52:49 閱讀更多