Pipeline rapid
Înainte să procesezi date mai complexe, managerul tău ar dori să vadă un exemplu simplu de pipeline care include pașii de bază. În acest exemplu, vei încărca un fișier de date, vei filtra câteva rânduri, vei adăuga o coloană cu ID și vei scrie rezultatul ca date JSON.
Contextul spark este deja definit, iar biblioteca pyspark.sql.functions este importată cu aliasul F, conform convenției uzuale.
Acest exercițiu face parte din cursul
Curățarea datelor cu PySpark
Instrucțiuni pentru exercițiu
- Importă fișierul
2015-departures.csv.gzîntr-un DataFrame. Reține că antetul este deja definit. - Filtrează DataFrame-ul astfel încât să conțină doar zboruri cu o durată mai mare de 0 minute. Folosește indexul coloanei, nu numele acesteia (nu uita să folosești
.printSchema()pentru a vedea numele și ordinea coloanelor). - Adaugă o coloană cu ID.
- Scrie fișierul ca document JSON cu numele
output.json.
Exercițiu interactiv practic
Încearcă acest exercițiu completând acest cod de exemplu.
# Import the data to a DataFrame
departures_df = spark.____(____, header=____)
# Remove any duration of 0
departures_df = departures_df.____(____)
# Add an ID column
departures_df = departures_df.____('id', ____)
# Write the file out to JSON format
____.write.____(____, mode='overwrite')