パーティション分割と遅延処理

PySpark でデータをクレンジングする

Mike Metzger

Data Engineering Consultant

パーティション分割

  • DataFrame はパーティションに分割されます
  • パーティションサイズは可変
  • 各パーティションは独立に処理
PySpark でデータをクレンジングする

遅延処理

  • 変換は遅延評価です
    • .withColumn(...)
    • .select(...)
  • アクションを実行するまで実際には処理されません
    • .count()
    • .write(...)
  • 変換は最適化のため並べ替えられることがあります
  • 想定外の挙動を招く場合があります
PySpark でデータをクレンジングする

ID の追加

通常の ID フィールド:

  • リレーショナル DB で一般的
  • 多くは増加する連番の一意な整数
  • 並列化しにくい
id last name first name state
0 Smith John TX
1 Wilson A. IL
2 Adams Wendy OR
PySpark でデータをクレンジングする

単調増加 ID

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 は遅延評価であることに注意してください。

  • ときどき順序が前後します
  • 結合を行う場合、ID は結合後に付与されることがあります
  • 変換は必ずテストしましょう
PySpark でデータをクレンジングする

実践してみましょう!

PySpark でデータをクレンジングする

Preparing Video For Download...