【问题标题】:How to define a global read\write variables in Spark如何在 Spark 中定义全局读写变量
【发布时间】:2020-07-22 18:30:03
【问题描述】:

Spark 有broadcast 变量,它们是只读的,accumulator 变量,可以由节点更新,但不能读取。有没有办法 - 或解决方法 - 定义一个既可更新又可读取的变量?

对这种读\写全局变量的一个要求是实现缓存。当文件作为 rdd 加载和处理时,会执行计算。这些计算的结果——发生在多个并行运行的节点中——需要放入一个映射中,该映射具有正在处理的实体的一些属性作为关键。随着 rdd 中的后续实体被处理,缓存被查询。

Scala 确实有 ScalaCache,它是缓存实现的外观,例如 Google Guava。但是如何在 Spark 应用程序中包含和访问这样的缓存呢?

缓存可以定义为驱动程序应用程序中创建SparkContext 的变量。但是接下来会有两个问题:

  • 由于网络开销,性能可能会很差 在节点和驱动程序应用程序之间。
  • 据我了解,每个 rdd 都会传递一个变量的副本 (在这种情况下为缓存)当变量第一次被 函数传递给 rdd。每个 rdd 都有自己的副本,不能访问共享的全局变量。

实现和存储这种缓存的最佳方式是什么?

谢谢

【问题讨论】:

  • 如何在 Spark 中定义全局读写变量,例如用于定义缓存,如我的示例所示。
  • 感谢 Tzach - 将为该问题添加新评论

标签: apache-spark


【解决方案1】:

嗯,最好的办法就是不这样做。一般来说,Spark 处理模型不提供任何保证*关于

  • 在哪里,
  • 什么时候,
  • 按什么顺序(当然不包括由血统/DAG 定义的转换顺序)
  • 多少次

给定的一段代码被执行。此外,任何直接依赖于 Spark 架构的更新都不是粒度的。

这些属性使 Spark 具有可扩展性和弹性,但同时也使得保持共享可变状态非常难以实现,并且在大多数情况下完全无用。

如果你想要的只是一个简单的缓存,那么你有多种选择:

  • 使用Tzach ZoharCaching in Spark 中描述的方法之一
  • 使用本地缓存(每个 JVM 或执行程序线程)结合应用程序特定的分区来保持本地化
  • 与外部系统通信使用独立于 Spark 的节点本地缓存(例如用于 http 请求的 Nginx 代理)

如果应用程序需要更复杂的通信,您可以尝试不同的消息传递工具来保持同步状态,但通常它需要复杂且可能很脆弱的代码。


* 这在 Spark 2.4 中发生了部分变化,引入了屏障执行模式(SPARK-24795SPARK-24822)。

【讨论】:

  • Spark StreamingLinearAlgorithm 中有一个类,其中模型对象被更新并用于预测。这不符合对同一对象进行读写的示例。我不确定,如果你能解释一下。
  • 我也在这里问过一个与此相关的问题。 stackoverflow.com/questions/43114971/…
猜你喜欢
  • 1970-01-01
  • 2018-01-27
  • 2016-05-07
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2016-02-18
  • 2011-08-20
相关资源
最近更新 更多