Pipeline rapida
Prima di analizzare dati più complessi, la tua responsabile vuole vedere un esempio semplice di pipeline con i passaggi di base. In questo esempio importerai un file di dati, filtrerai alcune righe, aggiungerai una colonna ID e infine scriverai il risultato come dati JSON.
Il contesto spark è già definito e la libreria pyspark.sql.functions è importata con alias F, come da consuetudine.
Questo esercizio fa parte del corso
Pulizia dei dati con PySpark
Istruzioni dell'esercizio
- Importa il file
2015-departures.csv.gzin un DataFrame. Nota che l'intestazione è già definita. - Filtra il DataFrame per mantenere solo i voli con durata superiore a 0 minuti. Usa l’indice della colonna, non il nome della colonna (ricorda di usare
.printSchema()per vedere nomi e ordine delle colonne). - Aggiungi una colonna ID.
- Scrivi il file come documento JSON chiamato
output.json.
esercizio interattivo pratico
Prova questo esercizio completando questo codice di esempio.
# 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')