我认为这个问题最好表述为:
我们什么时候需要调用缓存或持久化 RDD?
Spark 进程是惰性的,也就是说,除非需要,否则不会发生任何事情。
为了快速回答这个问题,在发出val textFile = sc.textFile("/user/emp.txt") 之后,数据没有任何变化,只构造了一个HadoopRDD,使用文件作为源。
假设我们对数据进行了一些转换:
val wordsRDD = textFile.flatMap(line => line.split("\\W"))
同样,数据没有任何反应。现在有一个新的 RDD wordsRDD 包含对 testFile 的引用和一个在需要时应用的函数。
只有当一个动作被一个RDD调用时,比如wordsRDD.count,RDD链,称为lineage才会被执行。也就是说,数据被划分为分区,将由 Spark 集群的 executors 加载,应用 flatMap 函数并计算结果。
在线性谱系上,例如本例中的谱系,不需要cache()。数据将被加载到执行程序,所有转换都将被应用,最后count 将被计算,全部在内存中 - 如果数据适合内存。
cache 在 RDD 的血统分支出来时很有用。假设您要将上一个示例中的单词过滤为正负单词的计数。你可以这样做:
val positiveWordsCount = wordsRDD.filter(word => isPositive(word)).count()
val negativeWordsCount = wordsRDD.filter(word => isNegative(word)).count()
在这里,每个分支都会重新加载数据。添加显式的cache 语句将确保之前完成的处理被保留和重用。作业将如下所示:
val textFile = sc.textFile("/user/emp.txt")
val wordsRDD = textFile.flatMap(line => line.split("\\W"))
wordsRDD.cache()
val positiveWordsCount = wordsRDD.filter(word => isPositive(word)).count()
val negativeWordsCount = wordsRDD.filter(word => isNegative(word)).count()
因此,cache 被称为“打破血统”,因为它创建了一个检查点,可重复用于进一步处理。
经验法则:当您的 RDD 的血统分支出或当一个 RDD 被多次使用时(例如在循环中),请使用 cache。