【问题标题】:Spark Streaming: how to propagate updates to a Broadcast variable to the whole cluster?Spark Streaming:如何将广播变量的更新传播到整个集群?
【发布时间】:2015-08-06 20:53:25
【问题描述】:

我在 Spark 驱动程序中有一个模块正在侦听 Kafka 队列,并且根据队列的内容,我需要修改广播变量(或闭包)的内容。在此示例中,这可能是一个字符串。

例如,如果字符串“change”到达队列,我需要更新每个节点中的广播变量。

我希望看到一种干净且高效的模式来执行此操作,或者至少接收关于我可以在哪里找到一些材料以更好地了解如何在 Spark 集群中传播修改的输入。

【问题讨论】:

  • 您能否提供一些示例代码来说明您想要做什么?
  • 特别是“我在 Spark 驱动程序中有一个模块正在监听 Kafka 队列”令人费解。有关详细信息,请参阅@Bacon 答案的讨论。

标签: scala apache-spark spark-streaming


【解决方案1】:

广播变量确实是使用点对点协议将变量或整个闭包传播到 spark 集群。

来自Learning Spark book

广播变量是 只是 spark.broadcast.Broadcast[T] 类型的对象,它包装了一个值 类型 T。我们可以通过在我们的广播对象上调用 value 来访问这个值 任务。该值仅发送到每个节点一次,使用高效的类似 BitTorrent 沟通机制。

影响性能的是您正在使用的序列化方法(例如:Kryo、自定义方法...):

书中有一个例子:

示例 6-8。在 Scala 中使用广播值查找国家/地区

// Look up the countries for each call sign for the
// contactCounts RDD. We load an array of call sign
// prefixes to country code to support this lookup.

val signPrefixes = sc.broadcast(loadCallSignTable())

val countryContactCounts = contactCounts.map {
    case (sign, count) =>
        val country = lookupInArray(sign, signPrefixes.value) (country, count)
    }.reduceByKey((x, y) => x + y)

countryContactCounts.saveAsTextFile(outputDir + "/countries.txt")

如这些示例所示,使用广播变量的过程很简单: 1. 通过在类型对象上调用 SparkContext.broadcast 创建一个 Broadcast[T] T. 只要它也是可序列化的,任何类型都可以工作。 2. 使用 value 属性(或 Java 中的 value() 方法)访问其值。 3. 变量只会被发送到每个节点一次,并且应该被视为读取- 仅(更新不会传播到其他节点)。

【讨论】:

  • “广播变量确实是使用点对点协议将变量或整个闭包传播到 spark 集群。”好的,这是我不确定的事情。我将在 sn-p 中尝试它,看看它的行为如何。谢谢。
  • AFAIK 广播变量只能在驱动程序上声明(如示例所示),如果需要根据在一个执行程序处接收到的内容更新某些值,则此方法将不起作用。
  • OP 想要向执行者广播一个在驱动程序上接收到的值。
  • @Bacon 对我来说,这个问题还不清楚。给定“例如,如果字符串“change”到达队列,我需要更新每个节点中的广播变量。” -- 数据到达执行者。只有collect()(或类似的)才能将其带给司机,这可能是也可能不是和选项。
  • OP 声明驱动程序正在侦听 Kafka 队列,并且这些消息应该在执行程序上转发。
猜你喜欢
  • 2016-08-18
  • 1970-01-01
  • 2015-08-26
  • 2016-11-18
  • 2016-05-26
  • 2015-05-06
  • 1970-01-01
  • 2015-04-18
  • 1970-01-01
相关资源
最近更新 更多