PySpark 的彈性分散式資料集

PySpark 入門

Benjamin Schmidt

Data Engineer

PySpark 中的平行化是什麼?

  • 自動在叢集中多個節點分散資料與計算
  • 在多個節點上分散處理大型資料集
  • Worker 節點並行處理,作業結束時再彙整
  • 規模化後更快(可達 GB 甚至 TB)

平行化

PySpark 入門

認識 RDD

RDDs,即 Resilient Distributed Datasets:

  • 在叢集間分散的資料集合,節點故障可自動復原
  • 適合大規模資料
  • 不可變,可用 map()filter() 等轉換;用 collect() 取回結果或用 paralelize() 建立 RDDs
PySpark 入門

建立 RDD

# Initialize a Spark session
from pyspark.sql import SparkSession

spark = SparkSession.builder.appName("RDDExample").getOrCreate()

# Create a DataFrame from a csv census_df = spark.read.csv("/census.csv")
# Convert DataFrame to RDD census_rdd = census_df.rdd
# Show the RDD's contents using collect() census_rdd.collect()
PySpark 入門

示範 collect

# Collect the entire DataFrame into a local Python list of Row objects
data_collected = df.collect()

# Print the collected data
for row in data_collected:
    print(row)
```    
PySpark 入門

RDD 與 DataFrame 比較

DataFrame

  • 高階:著重易用與最佳化
  • 類 SQL 操作:用類 SQL 查詢,少量程式碼即可完成複雜運算
  • 綱要資訊:含欄與型別,類似 SQL 資料表

RDD

  • 低階:更彈性,但複雜操作需要較多程式碼
  • 型別保留:保留資料型別,但沒有 DataFrame 的最佳化
  • 無綱要:處理結構化資料(如 SQL、關聯式資料)較不便
  • 易於大規模擴充
  • 相較 DataFrame 冗長,分析表現較差
PySpark 入門

常用函式與方法

  • map():將函式(包含你自寫如 lambda)套用到資料集,例如: rdd.map(map_function)
  • collect():從叢集中收集資料,例如: rdd.collect()
PySpark 入門

一起來練習吧!

PySpark 入門

Preparing Video For Download...