Examinarea rândurilor invalide
Ai filtrat cu succes rândurile folosind un join, dar uneori vrei să examinezi datele care sunt invalide. Aceste date pot fi stocate pentru procesare ulterioară sau pentru depanarea surselor de date.
Vrei să găsești diferența dintre două DataFrame-uri și să stochezi rândurile invalide.
Obiectul spark este definit, iar pyspark.sql.functions este importat ca F. DataFrame-ul original split_df și DataFrame-ul rezultat în urma join-ului joined_df sunt disponibile în stările lor anterioare.
Acest exercițiu face parte din cursul
Curățarea datelor cu PySpark
Instrucțiuni pentru exercițiu
- Determină numărul de rânduri pentru fiecare DataFrame.
- Creează un DataFrame care conține doar rândurile invalide.
- Validează că numărul de rânduri din noul DataFrame este cel așteptat.
- Determină numărul de rânduri de tip folder distincte care au fost eliminate.
Exercițiu interactiv practic
Încearcă acest exercițiu completând acest cod de exemplu.
# Determine the row counts for each DataFrame
split_count = ____
joined_count = ____
# Create a DataFrame containing the invalid rows
invalid_df = split_df.____(____(joined_df), '____', '____')
# Validate the count of the new DataFrame is as expected
invalid_count = ____
print(" split_df:\t%d\n joined_df:\t%d\n invalid_df: \t%d" % (split_count, joined_count, invalid_count))
# Determine the number of distinct folder rows removed
invalid_folder_count = invalid_df.____('____').____.____
print("%d distinct invalid folders found" % invalid_folder_count)