Introduktion till datapipelines

Datarensning med PySpark

Mike Metzger

Data Engineering Consultant

Vad är en datapipeline?

  • En serie steg som bearbetar data från källa till slutresultat
  • Kan bestå av valfritt antal steg eller komponenter
  • Kan sträcka sig över många system
  • Vi fokuserar på datapipelines i Spark
Datarensning med PySpark

Hur ser en datapipeline ut?

  • Indata
    • CSV, JSON, webbtjänster, databaser
  • Transformationer
    • withColumn(), .filter(), .drop()
  • Utdata
    • CSV, Parquet, databas
  • Validering
  • Analys
Datarensning med PySpark

Om pipelines

  • Inte formellt definierad i Spark
  • Vanligtvis all normal Spark-kod som krävs för uppgiften
    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')
    
Datarensning med PySpark

Nu kör vi en övning!

Datarensning med PySpark

Preparing Video For Download...