什麼是 Spark DataFrame API?

Apache Spark 是分散式處理引擎,專為大規模處理結構化、半結構化和非結構化資料而設計。開發人員可使用 Spark API 與引擎互動,並定義分散式資料作業。


現代 Spark 開發的核心是 Spark DataFrame API。DataFrame 會將分散式資料整理及處理成資料欄並設定名稱,類似關聯式資料庫中的資料表或試算表。提供結構化的宣告式介面,可分析龐大資料集,同時在幕後自動最佳化實體執行作業。

Spark API 架構與執行

如要瞭解 Spark API 如何在分散式網路中執行程式碼,請務必瞭解其執行階段架構的基本元件:


  • 驅動程式:Spark 應用程式的中央協調器。它會讀取程式碼、將宣告式作業轉譯為邏輯執行計畫,並在 worker 節點間排程工作。
  • 叢集管理工具:資源分配工具 (例如 Kubernetes 或獨立資源管理工具),負責在整個叢集中佈建 CPU、記憶體和網路資源。
  • 執行器:在叢集節點上執行的 worker 執行個體。執行器會從驅動程式接收工作指令,在本機記憶體中執行資料處理作業,並傳回結果或將結果寫入儲存空間。
  • DAG、階段和工作:Spark 不會逐行執行資料轉換。相反地,驅動程式會將您的 DataFrame 程式碼編譯成由實體運算子組成的有向無環圖 (DAG)。驅動程式會將這些運算子分組為廣泛的執行階段 (以 data shuffling 邊界劃分),並將這些階段細分為個別工作,分配給執行器平行執行。

Spark 彈性分散式資料集 (RDD) 和 DataFrame 有什麼差別?

設計分散式資料管道時,開發人員可選擇兩種主要的程式設計介面:


  • Spark RDD (彈性分散式資料集):原始的低階 Spark API。代表不可變、容錯的 JVM 物件集合,分散在叢集中。由於 RDD 包含任意 Java 物件,因此執行引擎無法檢查其內部結構。這迫使開發人員手動編寫及調整低階轉換邏輯,往往導致垃圾回收負擔過重,以及序列化速度緩慢。
  • Spark DataFrame:現代化結構化資料處理標準。DataFrame 會根據嚴格的結構定義,將資料整理成資料列和資料欄。由於執行引擎瞭解資料集的資料類型和結構,因此可以在執行程式碼之前自動最佳化。

選擇 DataFrame 而非彈性分散式資料集 (RDD) 的理由

幾乎所有資料工程和資料科學用途,都適合使用 DataFrame API。彈性分散式資料集 (RDD) 僅適用於必須使用自訂 Java 物件,來操控原始非結構化資料 (例如二進位檔案或媒體串流) 的情況。DataFrame 需要的程式碼較少、自動執行速度更快,並利用堆積外二進位檔記憶體管理,略過 JVM 垃圾回收瓶頸。

資料團隊如何使用 Spark API

Spark DataFrame API 可簡化需要密集運算的任務,例如處理、清理及準備大量資料。


  • 資料工程師:工程師使用 DataFrame API 建構具備韌性和容錯能力的資料管道。他們會擷取原始資料、套用結構定義、執行結構化轉換,並將符合規範的資料表寫入資料湖泊或雲端資料倉儲。宣告式語法讓他們能每天處理數 TB 的資料,同時盡可能降低程式碼複雜度和維護負擔。
  • 資料科學家:資料科學家使用 DataFrame API 探索及準備龐大的資料集,這些資料集的大小超過單一機器的記憶體限制。只要將資料分區分散到各個 worker 節點,就能大規模執行探索性資料分析、清理空值及進行特徵工程,更快取得洞察資訊。

以 DataFrame API 為基礎建構的核心程式庫

Spark DataFrame API 是 Spark 進階處理程式庫的程式設計基礎:


  • Spark SQL:這個模組可讓您直接對 Spark 資料集執行 ANSI SQL 查詢。您可以將 DataFrame 當做暫時 view 查詢,完美結合宣告式 Python 或 Scala 程式碼與標準 SQL 查詢。
  • 結構化串流:這個引擎會處理連續的即時資料串流。這項功能使用與靜態批次處理相同的 DataFrame API 指令,會自動處理微批次、容錯和狀態管理。
  • MLlib:Spark 的分散式機器學習程式庫,使用 DataFrame 管理及準備訓練資料集,提供內建的可擴充演算法,用於分類、迴歸、分群和協同過濾。

Spark DataFrame API 的優點

宣告式方法最佳化

DataFrame 使用內建的查詢最佳化器,稱為 Catalyst 最佳化器。編寫 DataFrame 程式碼時,最佳化器會自動重寫實體執行計畫,套用篩選器下推、投影剪枝和 join 重新排序,盡可能加快執行速度,無需手動調整。

語言彈性

無論團隊使用 Python (PySpark)、Scala、Java 或 R 編寫程式碼,DataFrame API 都能確保效能一致。基礎執行計畫會在同一個獨立於 JVM 的執行層中編譯及最佳化。

有效率的記憶體管理

DataFrame 採用 Tungsten 執行引擎,以高度壓縮的堆積外二進位檔格式儲存資料。這可消除 JVM 物件建立負擔,並防止因凍結執行器執行緒而暫停垃圾回收。

如何在 Google Cloud 使用 Spark DataFrame API

傳統上,執行開放原始碼 Spark 需要手動設定叢集、管理軟體版本,以及進行複雜的網路對接。Google Cloud 提供 Managed Service for Apache Spark,將分散式執行環境轉化為全代管、企業級的資料平台,簡化這項作業。


資料團隊可使用下列功能,在 Google Cloud 執行 Spark DataFrame 工作負載:


  • 無伺服器和代管叢集部署:Google Cloud 讓您選擇符合營運需求的執行作業模型。您可以在無伺服器模式下執行 DataFrame 程式碼,直接提交批次工作,並在資源自動調度時,僅依執行階段的實際秒數付費。或者,您也可以部署高度可自訂的永久性代管叢集,用於持續性工作負載。
  • 最佳化 BigQuery 連線能力:spark-bigquery-connector 會直接以原生 Apache Arrow 格式使用 BigQuery 資料,藉此略過 JVM 資料列導向轉換。資料科學家也能使用整合式筆記本介面,在同一個受管理資料集上執行 PySpark DataFrame 程式碼和 SQL 查詢,不必切換環境。
  • 向量化原生 C++ 執行作業:在 Google Cloud 執行 Spark 工作負載時,您可以為無伺服器批次或代管叢集啟用 Lightning Engine。這個原生 C++ 查詢執行引擎會使用 Velox 和 Gluten,直接將實體查詢計畫編譯為原生指令。這項服務會略過 JVM Volcano 疊代器模式和垃圾回收瓶頸,因此可將 DataFrame 和 Spark SQL 工作負載的速度提升高達 4.9 倍,而且完全不需要變更程式碼。


展開下一步行動

運用價值 $300 美元的免費抵免額和超過 20 項一律免費的產品,開始在 Google Cloud 中建構產品與服務。