SQL și Parquet
Fișierele Parquet sunt ideale ca sursă de date pentru interogările SQL în Spark. Deși poți rula aceleași interogări direct prin funcțiile Python ale Spark, uneori e mai simplu să folosești interogări SQL alături de opțiunile Python.
În acest exemplu, vom citi fișierul Parquet creat în exercițiul anterior și îl vom înregistra ca tabel SQL. După înregistrare, vom rula o interogare rapidă asupra tabelului (adică a fișierului Parquet).
Obiectul spark și fișierul AA_DFW_ALL.parquet sunt disponibile automat.
Acest exercițiu face parte din cursul
Curățarea datelor cu PySpark
Instrucțiuni pentru exercițiu
- Importă fișierul
AA_DFW_ALL.parquetînflights_df. - Folosește metoda
createOrReplaceTempViewpentru a crea un alias pentru tabelulflights. - Rulează o interogare Spark SQL asupra tabelului
flights.
Exercițiu interactiv practic
Încearcă acest exercițiu completând acest cod de exemplu.
# Read the Parquet file into flights_df
flights_df = spark.read.____(____)
# Register the temp table
flights_df.____('flights')
# Run a SQL query of the average flight duration
avg_duration = spark.____('SELECT avg(flight_duration) from flights').collect()[0]
print('The average flight time is: %d' % avg_duration)