缓存

使用 PySpark 进行数据清洗

Mike Metzger

Data Engineering Consultant

什么是缓存?

Spark 中的缓存:

  • 将 DataFrame 存在内存或磁盘
  • 加速后续转换/动作
  • 降低资源消耗
使用 PySpark 进行数据清洗

缓存的缺点

  • 特大数据集可能装不下内存
  • 本地磁盘缓存未必提升性能
  • 缓存对象可能不可用
使用 PySpark 进行数据清洗

缓存技巧

开发 Spark 任务时:

  • 仅在需要时缓存
  • 在不同节点尝试缓存并评估是否提速
  • 优先缓存到内存和快速 SSD/NVMe
  • 需要时缓存到慢速本地磁盘
  • 使用中间文件!
  • 完成后停止缓存对象
使用 PySpark 进行数据清洗

实现缓存

在动作前对 DataFrame 调用 .cache()

voter_df = spark.read.csv('voter_data.txt.gz')
voter_df.cache().count()
voter_df = voter_df.withColumn('ID', monotonically_increasing_id())
voter_df = voter_df.cache()
voter_df.show()
使用 PySpark 进行数据清洗

更多缓存操作

.is_cached 查看缓存状态

print(voter_df.is_cached)
True

完成后对 DataFrame 调用 .unpersist()

voter_df.unpersist()
使用 PySpark 进行数据清洗

Passons à la pratique !

使用 PySpark 进行数据清洗

Preparing Video For Download...