Validează rândurile printr-un join
Un alt mod de a filtra datele este folosirea joinurilor pentru a elimina înregistrările invalide. Va trebui să verifici că numele folderelor corespund celor așteptate, pe baza unui DataFrame numit valid_folders_df. DataFrame-ul split_df este cel cu care ai lucrat anterior, conținând un grup de coloane divizate.
Obiectul spark este disponibil, iar pyspark.sql.functions este importat ca F.
Acest exercițiu face parte din cursul
Curățarea datelor cu PySpark
Instrucțiuni pentru exercițiu
- Redenumește coloana
_c0înfolderîn DataFrame-ulvalid_folders_df. - Numără rândurile din
split_df. - Realizează un join între cele două DataFrame-uri pe baza numelui folderului și numește DataFrame-ul rezultat
joined_df. Asigură-te că aplici broadcast pe DataFrame-ul mai mic. - Verifică numărul de rânduri rămase în DataFrame și compară rezultatele.
Exercițiu interactiv practic
Încearcă acest exercițiu completând acest cod de exemplu.
# Rename the column in valid_folders_df
valid_folders_df = ____
# Count the number of rows in split_df
split_count = ____
# Join the DataFrames
joined_df = split_df.____(____(valid_folders_df), "folder")
# Compare the number of rows remaining
joined_count = ____
print("Before: %d\nAfter: %d" % (split_count, joined_count))