Apache Spark中仅复用两次的Dataset是否需要缓存?
Apache Spark Dataset缓存必要性与性能提升分析
问题场景
我正在使用Apache Spark开发,编写了如下代码:
Dataset<Row> tradesDataset = sparkSession .sql("select * from a_table") .cache(); // <-- do I need caching here? long countOfDistinctUitIdsInTradeAgreements = tradesDataset .select(tradesDataset.col("uitid")) .distinct() .count(); long countOfDistinctUitIdsInTradeAgreementsForTradeDate = tradesDataset .filter(tradesDataset.col("TRADE_DATE").equalTo(processingDate)) .count();
基于该Dataset执行两次不同的统计操作,咨询:是否需要缓存该查询生成的Dataset?此举能否带来性能提升?
回答
需要缓存,且能显著提升性能
Spark的Dataset采用懒执行机制,若不添加cache(),两次统计操作会各自触发一次select * from a_table的全表扫描——相当于重复读取并计算全表数据,IO和计算成本直接翻倍。添加缓存后,全表扫描仅执行一次,结果会被存储到Spark集群的内存(或磁盘,取决于默认的MEMORY_AND_DISK存储级别),后续两次统计直接基于缓存数据计算,避免了重复的全表扫描开销。优化细节建议
- 按需调整缓存级别:默认
cache()对应MEMORY_AND_DISK,如果数据量较小且集群内存充足,可改用persist(StorageLevel.MEMORY_ONLY)以获得更快的读取速度;若数据量较大,MEMORY_AND_DISK是更稳妥的选择,内存放不下的部分会自动写入磁盘。 - 缩小缓存数据范围:观察两次统计操作仅用到
uitid和TRADE_DATE字段,可先对原查询做字段裁剪后再缓存,减少缓存的数据体积,进一步提升效率:Dataset<Row> tradesDataset = sparkSession .sql("select uitid, TRADE_DATE from a_table") .cache(); - 及时释放缓存资源:如果后续代码不再使用
tradesDataset,记得调用tradesDataset.unpersist()手动释放缓存,避免长期占用集群内存资源。
内容的提问来源于stack exchange,提问作者Sergey Tsypanov
相关产品推荐
相关产品推荐

