資料管線導論

使用 PySpark 清理資料

Mike Metzger

Data Engineering Consultant

什麼是資料管線?

  • 一組步驟,將資料由來源處理到最終輸出
  • 可包含任意數量的步驟或元件
  • 可跨越多個系統
  • 本章聚焦於 Spark 中的資料管線
使用 PySpark 清理資料

資料管線長什麼樣?

  • 輸入
    • CSV、JSON、網路服務、資料庫
  • 轉換
    • withColumn().filter().drop()
  • 輸出
    • CSV、Parquet、資料庫
  • 驗證
  • 分析
使用 PySpark 清理資料

管線細節

  • 在 Spark 中沒有正式定義
  • 通常就是完成任務所需的一般 Spark 程式碼
    schema = StructType([
    StructField('name', StringType(), False),
    StructField('age', StringType(), False)
    ])
    df = spark.read.format('csv').load('datafile').schema(schema)
    df = df.withColumn('id', monotonically_increasing_id())
    ...
    df.write.parquet('outdata.parquet')
    df.write.json('outdata.json')
    
使用 PySpark 清理資料

一起來練習吧!

使用 PySpark 清理資料

Preparing Video For Download...