Szybki potok danych
Zanim przejdziesz do przetwarzania bardziej złożonych danych, twój menedżer chciałby zobaczyć prosty przykład potoku danych obejmujący podstawowe kroki. W tym ćwiczeniu wczytasz plik z danymi, przefiltruj kilka wierszy, dodaj kolumnę z identyfikatorem, a następnie zapisz wynik jako dane JSON.
Kontekst spark jest już zdefiniowany, a biblioteka pyspark.sql.functions jest zaimportowana z aliasem F, zgodnie z przyjętą konwencją.
To ćwiczenie jest częścią kursu
Czyszczenie danych w PySpark
Instrukcje do ćwiczenia
- Wczytaj plik
2015-departures.csv.gzdo DataFrame. Nagłówek jest już zdefiniowany. - Przefiltruj DataFrame tak, aby zawierał tylko loty trwające ponad 0 minut. Użyj indeksu kolumny, a nie jej nazwy (pamiętaj, że możesz użyć
.printSchema(), aby sprawdzić nazwy i kolejność kolumn). - Dodaj kolumnę z identyfikatorem.
- Zapisz wynik jako dokument JSON o nazwie
output.json.
Interaktywne ćwiczenie praktyczne
Spróbuj tego ćwiczenia, uzupełniając ten przykładowy kod.
# 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')