【问题标题】:spark - RDD process twice after persistspark - 持久化后两次 RDD 进程
【发布时间】:2018-11-22 10:26:53
【问题描述】:

我创建了一个 RDD 并从 origin 创建了另一个 RDD,如下所示。

val RDD2 = RDD1.map({
  println("RDD1")
  ....
}).persist(StorageLevel.MEMORY_AND_DISK)

RDD2.foreach({
  println("RDD2")
  ...
})
...so on..

我预计RDD1的进程只会执行一次,因为RDD1是通过persist方法保存在内存或磁盘上的。

但不知何故,“RDD1”打印在“RDD2”之后,如下所示。

RDD1
RDD1
RDD1
RDD1
RDD2
RDD2
RDD2
RDD2
RDD2
RDD1 -- repeat RDD1 process. WHY? 
RDD1
RDD1
RDD1
RDD2
RDD2
RDD2
RDD2
RDD2

【问题讨论】:

  • 当第一个“RDD1”全部打印出来时,我可以保证RDD1的处理已经完成。它做了两次相同的工作。
  • 我猜你最后做了一些collect 动作? spark.apache.org/docs/latest/api/java/org/apache/spark/rdd/… 也返回一个 RDD。
  • @davidshen84 是的,我在代码末尾做了 collectAsMap()
  • 对于 RDD1RDD2?请记住 Spark 中的两个主要概念,即转换和动作。在您采取行动之前,转换不会对 RDD 产生任何影响。

标签: apache-spark


【解决方案1】:

这是 spark 的预期行为。像大多数spark中坚持的操作也是懒惰的操作。因此,即使您为第一个 RDD 添加了持久化,spark 也不会缓存数据,除非您在持久化操作之后添加任何操作。 map操作在spark中不是action,也是惰性的。

强制缓存的方法是在RDD2的persist之后添加count动作

val RDD2 = RDD1.map({
   println("RDD1")
   ....
}).persist(StorageLevel.MEMORY_AND_DISK)

RDD2.count // Forces the caching 

现在,如果您执行任何其他操作,RDD2 将不会被重新计算

【讨论】:

  • 我们可以使用任何配置选项来强制缓存吗?即不调用操作方法
  • 不,这是 spark 设计的。您必须通过调用 Rdd 上的操作来强制缓存
  • 来自 Spark RDD 文档网页:“默认情况下,每个转换后的 RDD 可能会在您每次对其运行操作时重新计算。但是,您也可以使用持久化(或缓存)将 RDD 持久化到内存中方法,在这种情况下,Spark 会将元素保留在集群上,以便下次查询时更快地访问”
  • 这有点令人困惑,所以网页说persist使RDD被缓存以备将来使用,而您建议它是一个惰性操作
  • 它会保留在缓存中,但不会立即保留。它仅在对此调用操作后才会持续存在。如需参考,请参阅本书jaceklaskowski.gitbooks.io/mastering-spark-sql/…。这明确提到了火花缓存是惰性操作。
猜你喜欢
  • 2016-01-16
  • 2015-09-13
  • 2016-07-17
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多