【发布时间】:2017-11-22 01:51:01
【问题描述】:
在 CouchDB 和 Incoop 等系统设计中,有一个称为“增量 MapReduce”的概念,其中保存了以前执行 MapReduce 算法的结果,并用于跳过未更改的输入数据部分。
假设我有 100 万行分为 20 个分区。如果我对这些数据运行一个简单的 MapReduce,我可以缓存/存储减少每个单独分区的结果,然后再将它们组合并再次减少以产生最终结果。如果我只更改第 19 个分区中的数据,那么我只需要对数据的更改部分运行 map & reduce 步骤,然后将新结果与未更改分区中保存的 reduce 结果相结合,以获得更新的结果。使用这种捕获方式,我可以跳过几乎 95% 的工作,以便在这个假设的数据集上重新运行 MapReduce 作业。
有没有什么好的方法可以将此模式应用于 Spark?我知道我可以编写自己的工具来将输入数据拆分为分区,检查我之前是否已经处理过这些分区,如果有的话从缓存中加载它们,然后运行最终的 reduce 以将所有分区连接在一起。但是,我怀疑有一种更简单的方法可以解决这个问题。
我已经在 Spark Streaming 中尝试过检查点,它能够在重新启动之间存储结果,这几乎是我想要的,但我想在流作业之外执行此操作。
RDD 缓存/持久化/检查点几乎看起来像是我可以构建的东西——它可以很容易地保留中间计算并在以后引用它们,但我认为一旦 SparkContext 停止,缓存的 RDD 总是会被删除,即使它们'被持久化到磁盘。因此缓存不适用于在重新启动之间存储结果。另外,我不确定在启动新 SparkContext 时是否应该/如何加载检查点 RDD……它们似乎存储在特定于 SparkContext 单个实例的 UUID in the checkpoint directory 下。
【问题讨论】:
标签: apache-spark