检查无效行
您已经通过连接成功过滤掉了这些行,但有时您还希望检查哪些数据是无效的。无效数据可以被保存起来,以便后续处理或用于排查数据源问题。
现在,您想找出两个 DataFrame 之间的差异,并把无效的行保存下来。
spark 对象已定义,且已将 pyspark.sql.functions 以 F 导入。原始 DataFrame split_df 和连接后的 DataFrame joined_df 已按之前的状态可用。
本练习是课程的一部分
使用 PySpark 进行数据清洗
练习说明
- 分别计算每个 DataFrame 的行数。
- 创建仅包含无效行的 DataFrame。
- 验证新 DataFrame 的行数是否符合预期。
- 计算被移除的不同文件夹行数。
交互式实操练习
通过完成这段示例代码来试试这个练习。
# 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)