Rychlý pipeline
Než se pustíš do zpracování složitějších dat, tvůj manažer by rád viděl jednoduchý příklad pipeline zahrnující základní kroky. V tomto cvičení načteš datový soubor, vyfiltrujemeš několik řádků, přidáš sloupec s ID a výsledek zapíšeš jako JSON.
Kontext spark je již definován a knihovna pyspark.sql.functions je podle zvyklostí dostupná pod aliasem F.
Toto cvičení je součástí kurzu
Cleaning Data with PySpark
Pokyny k cvičení
- Importuj soubor
2015-departures.csv.gzdo DataFrame. Hlavička je již definována. - Vyfiltruj DataFrame tak, aby obsahoval pouze lety s dobou trvání delší než 0 minut. Použij index sloupce, nikoli jeho název (pro zobrazení názvů a pořadí sloupců použij
.printSchema()). - Přidej sloupec s ID.
- Zapiš soubor jako JSON dokument s názvem
output.json.
Interaktivní cvičení na vyzkoušení si v praxi
Vyzkoušejte si toto cvičení dokončením tohoto ukázkového kódu.
# 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')