Snabb pipeline
Innan du analyserar mer komplex data vill din chef se ett enkelt exempel på en pipeline med de grundläggande stegen. I det här exemplet ska du läsa in en datafil, filtrera bort några rader, lägga till en ID-kolumn och sedan skriva ut resultatet som JSON-data.
spark-kontexten är definierad och biblioteket pyspark.sql.functions är importerat med aliaset F, vilket är brukligt.
Den här övningen är en del av kursen
Datarensning med PySpark
Övningsinstruktioner
- Importera filen
2015-departures.csv.gztill en DataFrame. Observera att rubriken redan är definierad. - Filtrera DataFrame så att den endast innehåller flygningar med en varaktighet på mer än 0 minuter. Använd kolumnens index, inte kolumnnamnet (kör
.printSchema()för att se kolumnnamn och ordning). - Lägg till en ID-kolumn.
- Skriv ut filen som ett JSON-dokument med namnet
output.json.
Interaktiv övning med praktiskt arbete
Testa den här övningen genom att slutföra den här exempelkoden.
# 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')