请说明在 PySpark 中 cache() 与 persist() 的用途与实现机制,并解释它们如何帮助提升计算性能。
考察说明
考查对 PySpark 缓存机制的理解,包括 cache 与 persist 的区别及其对性能的影响。
回答思路
- 【回答框架 1】cache() 本质是 persist() 的特殊形式,默认使用 MEMORY_AND_DISK 存储级别,将 RDD 或 DataFrame 的计算结果存储于集群内存,必要时溢写到磁盘,从而避免重复计算同一数据集。
- 【回答框架 2】persist() 允许指定更细粒度的存储级别,如 MEMORY_ONLY、MEMORY_AND_DISK_SER、OFF_HEAP 等,以在内存占用、序列化开销和计算成本之间权衡。
- 【回答框架 3】性能提升源于避免重复执行相同的转换操作,尤其当同一数据集被多次用于不同动作(如多次 count 或迭代计算)时效果明显。同时需注意,缓存操作是惰性的,只有遇到动作算子时才会真正缓存数据。
- 【回答框架 4】选择缓存级别时需考虑数据大小、内存资源、序列化开销和重用频率。例如,内存充足时可选用 MEMORY_ONLY 避免序列化,数据较大时选 MEMORY_AND_DISK 防内存溢出。
- 【回答框架 5】使用一段时间后可用 unpersist() 释放缓存,避免占用集群资源。
- 【关键点 1】cache() 是 persist() 的默认级别 MEMORY_AND_DISK 的简写。
- 【关键点 2】缓存避免重复计算,但仅在遇到动作算子时才真正缓存。
- 【关键点 3】需根据数据大小和内存资源选择合适的存储级别。
- 【关键点 4】不再使用时调用 unpersist() 释放资源。
- 【易错点 1】缓存并不保证绝对避免重复计算,若数据源变化或节点故障可能导致重新计算。
- 【易错点 2】缓存占用内存,过度使用可能引发内存溢出,需合理设置存储级别。
- 【易错点 3】不要对同一数据集多次调用 cache(),只会在循环中产生冗余。