快取

使用 PySpark 清理資料

Mike Metzger

Data Engineering Consultant

什麼是快取?

Spark 的「快取」:

  • 將 DataFrame 存在記憶體或磁碟
  • 加速後續轉換/動作
  • 降低資源使用
使用 PySpark 清理資料

快取的缺點

  • 超大型資料集可能無法放入記憶體
  • 以本機磁碟快取未必能提昇效能
  • 已快取的物件可能不可用
使用 PySpark 清理資料

快取技巧

開發 Spark 工作時:

  • 需要時再快取
  • 嘗試在不同步驟快取 DataFrame,評估效能是否變好
  • 儘量用記憶體與快速 SSD/NVMe 儲存
  • 需要時再用較慢的本機磁碟
  • 善用中介檔!
  • 完成後停止快取物件
使用 PySpark 清理資料

實作快取

在執行 Action 前,先對 DataFrame 呼叫 .cache()

voter_df = spark.read.csv('voter_data.txt.gz')
voter_df.cache().count()
voter_df = voter_df.withColumn('ID', monotonically_increasing_id())
voter_df = voter_df.cache()
voter_df.show()
使用 PySpark 清理資料

更多快取操作

使用 .is_cached 檢查快取狀態

print(voter_df.is_cached)
True

完成後對 DataFrame 呼叫 .unpersist()

voter_df.unpersist()
使用 PySpark 清理資料

一起來練習吧!

使用 PySpark 清理資料

Preparing Video For Download...