大規模 PySpark

PySpark 入門

Benjamin Schmidt

Data Engineer

善用規模

  • PySpark 能有效處理 GB、TB 等級的資料
  • 使用 PySpark 的目標是追求速度與高效
  • 了解 PySpark 的執行機制可進一步提升效率
  • 使用 broadcast 來協調整個叢集
joined_df = large_df.join(broadcast(small_df), 
                          on="key_column", how="inner")
joined_df.show()
PySpark 入門

執行計畫

# 使用 explain() 檢視執行計畫
df.filter(df.Age > 40).select("Name").explain()
== Physical Plan ==
*(1) Filter (isnotnull(Age) AND (Age > 30))
+- Scan ExistingRDD[Name:String, Age:Int]
1 https://spark.apache.org/docs/latest/api/python/reference/pyspark.sql/api/pyspark.sql.DataFrame.explain.html
PySpark 入門

快取與永續化 DataFrame

  • 快取:將資料存於記憶體,適合較小資料集,存取更快
  • 永續化:以不同儲存層級保存,適合較大資料集
df = spark.read.csv("large_dataset.csv", header=True, inferSchema=True)

# 快取 DataFrame
df.cache()

# 對已快取的 DataFrame 執行多個操作 df.filter(df["column1"] > 50).show() df.groupBy("column2").count().show()
PySpark 入門

以不同儲存層級進行永續化

# 以儲存層級進行永續化
from pyspark import StorageLevel

df.persist(StorageLevel.MEMORY_AND_DISK)

# 進行轉換 result = df.groupBy("column3").agg({"column4": "sum"}) result.show() # 使用後解除永續化 df.unpersist()
PySpark 入門

最佳化 PySpark

  • 小區塊處理:資料越多,操作越慢;基於選擇性,優先用 map() 勝於 groupby()
  • 廣播連接:broadcast 會用滿整體算力,即使是小型資料集
  • 避免重複動作:對同一資料重複動作會耗時耗算力,沒有額外收益
PySpark 入門

一起來練習吧!

PySpark 入門

Preparing Video For Download...