分区与惰性计算

使用 PySpark 进行数据清洗

Mike Metzger

Data Engineering Consultant

分区

  • DataFrame 被拆分为分区
  • 分区大小可变
  • 各分区独立处理
使用 PySpark 进行数据清洗

惰性计算

  • 转换是惰性
    • .withColumn(...)
    • .select(...)
  • 只有执行 action 时才真正运行
    • .count()
    • .write(...)
  • 可重排转换以提升性能
  • 有时会产生意外行为
使用 PySpark 进行数据清洗

添加 ID

常规 ID 字段:

  • 关系型数据库中常见
  • 通常为递增、连续且唯一的整数
  • 并行性较差
id last name first name state
0 Smith John TX
1 Wilson A. IL
2 Adams Wendy OR
使用 PySpark 进行数据清洗

单调递增 ID

pyspark.sql.functions.monotonically_increasing_id()

  • 64 位整数,值单调递增,唯一
  • 不一定连续(可能有缺口)
  • 完全可并行
id last name first name state
0 Smith John TX
134520871 Wilson A. IL
675824594 Adams Wendy OR
使用 PySpark 进行数据清洗

注意事项

请记住,Spark 是"惰性"的!

  • 偶尔会出现乱序
  • 进行 join 时,ID 可能在 join 之后才分配
  • 请测试您的转换
使用 PySpark 进行数据清洗

让我们练习吧!

使用 PySpark 进行数据清洗

Preparing Video For Download...