提升匯入效能

使用 PySpark 清理資料

Mike Metzger

Data Engineering Consultant

Spark 叢集

「Spark 叢集」由兩種行程組成

  • Driver 行程
  • Worker 行程
使用 PySpark 清理資料

匯入效能

重要參數:

  • 物件數量(檔案、網路位置等)
    • 多個小物件優於單一大物件
    • 可用萬用字元匯入
      airport_df = spark.read.csv('airports-*.txt.gz')
      
  • 物件的一般大小
    • 若物件大小相近,Spark 表現更佳
使用 PySpark 清理資料

Schemas

良好定義的 schema 可大幅提升匯入效能

  • 避免重複讀取資料
  • 匯入時提供驗證
使用 PySpark 清理資料

如何分割物件

  • 使用作業系統工具/指令稿(split、cut、awk)
    split -l 10000 -d largefile chunk-
    
  • 使用自訂指令稿
  • 輸出為 Parquet
    df_csv = spark.read.csv('singlelargefile.csv')
    df_csv.write.parquet('data.parquet')
    df = spark.read.parquet('data.parquet')
    
使用 PySpark 清理資料

一起來練習吧!

使用 PySpark 清理資料

Preparing Video For Download...