データパイプライン入門

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

Mike Metzger

Data Engineering Consultant

データパイプラインとは?

  • ソースから最終出力までの処理手順の集合
  • 手順やコンポーネント数は任意
  • 複数システムにまたがる場合あり
  • 本講義では Spark 内のデータパイプラインに注目
PySpark でデータをクレンジングする

データパイプラインの構成

  • 入力
    • CSV、JSON、Web サービス、データベース
  • 変換
    • 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...