为什么已缓存的PySpark采样数据集执行count操作速度极慢?
核心原因分析
cache()是懒执行算子
Spark的缓存机制是懒加载的,你调用finaljoin_test_percent0001.cache()的时候,并不会立刻执行缓存操作,只是给该DataFrame打上了需要缓存的标记。只有第一次触发行动算子(比如你这里的count())时,才会完整执行「读取全量CSV→采样→缓存到指定存储介质→执行计数」的全链路流程,所以第一次count()的5分钟是包含了全量数据扫描、采样、缓存的总耗时,属于正常现象。你可以尝试对同一个DataFrame执行第二次count(),正常会降到秒级甚至毫秒级,就能验证缓存是否生效。大概率存在变量名笔误
你代码中缓存的变量名为finaljoin_test_percent0001,但执行计数的变量名为finaljoin_test_percent1,二者不是同一个对象。如果实际运行代码确实存在这个拼写错误,意味着你计数的是一个没有打缓存标记的DataFrame,相当于重新执行了一遍采样甚至全量扫描逻辑,自然不会用到缓存,耗时必然很高。
优化建议
- 先修正变量名拼写错误,确保所有后续操作都指向已经打了缓存标记的采样DataFrame
- 如果需要提前锁定缓存结果,可在
cache()之后主动调用一次轻量行动算子触发缓存加载,比如:finaljoin_test_percent0001 = finaljoin.sample(0.001).cache() # 主动触发缓存加载,后续操作直接读缓存 finaljoin_test_percent0001.count() - 超大规模CSV重复使用时建议提前转成Parquet等列式存储格式,读取速度相比CSV有数倍到数十倍提升,能大幅降低首次全量扫描的耗时
- 可通过
print(finaljoin_test_percent0001.storageLevel)查看当前缓存的存储级别,默认的MEMORY_ONLY在内存不足时会自动丢弃部分分区,后续用到时需要重算,可根据集群资源调整为MEMORY_AND_DISK避免重算开销。
内容的提问来源于stack exchange,提问作者Bratt Swan
相关产品推荐
相关产品推荐

