在即時風控、大促銷即時大屏、金融高頻交易與物聯網告警等高時效業務場景中,「資料一到達就必須在毫秒級完成計算」已成為現代企業的剛性需求。

在開源即時計算領域,Apache Flink 與 Apache Spark(Spark Streaming / Structured Streaming) 是佔據統治地位的兩大巨頭。

儘管兩者在現代版本中功能日益趨同,但在底層計算模型、狀態儲存架構與時間語義處理上,兩者存在著根本性的哲學分歧。

本文將帶你深入兩大引擎的底層核心,全面對比兩者的架構機制與選型維度。


1. 原生事件驅動(Flink)vs. 微批次(Spark)

兩大引擎最根本的差異在於「如何看待一條資料」:

Apache Flink 原生事件驅動 vs Spark 微批次架構對比圖展示 Flink 每到一筆 Event 立即觸發計算達亞毫秒級延遲,而 Spark 收集時間視窗微批次打包計算延遲約數百毫秒。✅ Apache Flink (原生事件驅動)每到 1 筆 Event ➔ 立即觸發計算 ➔ 更新 State立即向下游發送 (Stream-native)極致亞毫秒級延遲 (< 10ms)Spark Structured Streaming (微批次)時間視窗 (如 500ms) ➔ 收集 Events 打包 RDD啟動微批次排程計算 (Batch-native)延遲約 100ms ~ 500ms
評估維度Apache FlinkSpark Structured Streaming
計算本質事件驅動(Stream-native,批次被視為流的特例)微批次(Batch-native,流被視為連續的小批次)
端到端延遲極致低延遲(亞毫秒 ~ 數毫秒級)較低延遲(通常為 100ms ~ 500ms,Continuous Processing 模式有諸多算子限制)
吞吐量(Throughput)極高(透過 Credit-based Flow Control 優化)極致吞吐(批次向量化計算在超大數據量下極具優勢)
狀態維護(Stateful)原生有狀態運算(支援 TB 級 RocksDB 狀態後端)基於 Checkpoint 與 DeltaState 維護
生態整合與 Kafka、Pulsar、即時資料湖(Iceberg)深度整合與 Spark 離線批次、MLlib 機器學習、GraphX 完美統一

2. 時間語義與水位線機制(Watermark & Event Time)

在分散式網路環境中,由於網路延遲或設備重連,資料到達計算引擎的順序極易發生「亂序(Out-of-Order)」:

事件時間 (Event Time) 與處理時間 (Processing Time) 亂序對比圖展示事件按 12:00:01, 12:00:02, 12:00:03 產生,但引擎因延遲以 12:00:01, 12:00:05, 12:00:03 亂序接收。事件發生時間 (Event Time): 12:00:01 ──► 12:00:02 ──► 12:00:03 (物理真實時序)引擎接收時間 (Processing Time): 12:00:01 ──► 12:00:05 ──► 12:00:03 (網路延遲引發亂序!)

Watermark(水位線)如何解決亂序?

Flink 引入了 Watermark(水位線) 機制:

  • 水位線是一種特殊的控制事件,攜帶時間戳 t。
  • 當 Watermark(t) 到達時,代表引擎判定:所有 Event Time ≤ t 的資料均已被接收完畢,可以安全觸發視窗計算!
  • 延遲容忍(Allowed Lateness):若在 Watermark 推進後仍有極端遲到的資料到達,Flink 支援透過側輸出流(Side Output)將其捕獲至死信隊列,避免資料丟失。

在高階即時運算中(如滑動視窗聚合、雙流 JOIN、用戶連鎖行為識別),引擎必須在記憶體中維護龐大的中間狀態(State):

  • MemoryStateBackend:純 JVM 堆記憶體,適合輕量除錯。

  • FsStateBackend:狀態在記憶體,Checkpoint 存 S3/HDFS。

  • EmbeddedRocksDBStateBackend:狀態溢出至本地 NVMe SSD,突破 JVM GC 瓶頸!

  • RocksDB StateBackend 的威力:Flink 將狀態儲存在進程外的嵌入式 RocksDB(C++ 撰寫)中。當單節點狀態高達數百 GB 時,完全不受 Java GC 停頓影響,是支撐超大規模狀態運算的定海神針。


4. Exactly-Once 端到端精確一次語義保證

分散式節點可能隨時宕機崩潰,引擎如何保證「不重複計算、不遺漏計算」?

Flink Checkpoint Barrier 非同步屏障快照與 2PC Sink 架構圖展示 Barrier 流經 Source、Transform、Sink 算子觸發非同步快照寫入 S3/HDFS,成功時 2PC Sink 正式提交保證 Exactly-Once。Source (Kafka)注入 Barrier 1Barrier 1Transform 算子有狀態聚合計算Barrier 1Sink (2PC Commit)兩階段提交事務分散式持久化儲存 (AWS S3 / HDFS)快照 1: Kafka Offset + 算子 State + 2PC 預提交事務 ➔ 全員成功即標記 Checkpoint 成功!實現端到端 Exactly-Once 精確一次計算!
  1. Flink JobManager 定期向資料流注入 Checkpoint Barrier(屏障);
  2. 每個算子收到 Barrier 後,非同步將自身當前狀態寫入持久化儲存(S3 / HDFS);
  3. 當所有算子均回報快照完成,該 Checkpoint 被標記為成功。
  4. 配合 2PC(兩階段提交)Sink:下游輸出端(如 Kafka Producer 或 Iceberg)在 Checkpoint 成功時執行最終 Commit,實現真正的 端到端 Exactly-Once!

5. 選型決策指南

  • 選 Apache Flink:
    • 業務對延遲有嚴苛要求(< 100ms,如金融高頻風控、電競反作弊);
    • 需要複雜的即時狀態計算(超長視窗雙流 JOIN、CEP 複雜事件處理)。
  • 選 Spark Structured Streaming:
    • 業務對延遲容忍度在秒級以上(如 5 秒更新一次業務報表);
    • 團隊現有大數據技術棧高度依賴 Spark(希望離線批次與即時流共用同一套代碼與 ML Pipeline)。