分割與延遲處理

使用 PySpark 清理資料

Mike Metzger

Data Engineering Consultant

分割(Partitioning)

  • DataFrame 會被切成多個分割區(partitions)
  • 分割區大小可變
  • 每個分割區獨立處理
使用 PySpark 清理資料

延遲處理(Lazy processing)

  • 轉換是延遲
    • .withColumn(...)
    • .select(...)
  • 直到執行動作前實際不會計算
    • .count()
    • .write(...)
  • 轉換可重新排序以提升效能
  • 有時會造成意外行為
使用 PySpark 清理資料

加入 ID

一般 ID 欄位:

  • 關聯式資料庫常見
  • 多為遞增的整數,連續且唯一
  • 並行度不高
id last name first name state
0 Smith John TX
1 Wilson A. IL
2 Adams Wendy OR
使用 PySpark 清理資料

單調遞增 ID(Monotonically increasing IDs)

pyspark.sql.functions.monotonically_increasing_id()

  • 64 位整數,數值遞增且唯一
  • 不一定連續(可能有缺口)
  • 可完全並行
id last name first name state
0 Smith John TX
134520871 Wilson A. IL
675824594 Adams Wendy OR
使用 PySpark 清理資料

注意事項

記住,Spark 是「延遲」的!

  • 偶爾會亂序
  • 做 join 時,ID 可能在 join 之後才指派
  • 記得測試你的轉換
使用 PySpark 清理資料

一起來練習吧!

使用 PySpark 清理資料

Preparing Video For Download...