开始使用免费开始使用

检查无效行

您已经通过连接成功过滤掉了这些行,但有时您还希望检查哪些数据是无效的。无效数据可以被保存起来,以便后续处理或用于排查数据源问题。

现在,您想找出两个 DataFrame 之间的差异,并把无效的行保存下来。

spark 对象已定义,且已将 pyspark.sql.functionsF 导入。原始 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)
编辑并运行代码