【问题标题】:How to broadcast a variable in a Spark Streaming mapping function?如何在 Spark Streaming 映射函数中广播变量?
【发布时间】:2016-11-18 05:41:51
【问题描述】:

我知道常规:sc.broadcast(x)

但是,目前 Spark Streaming 不支持带有检查点的广播变量。

官方指南提供了解决方案:http://spark.apache.org/docs/latest/streaming-programming-guide.html#accumulators-and-broadcast-variables。但是,此方案只能用于 foreachRDD 函数。

现在我想在映射函数(如flatMapToPair)中使用需要以这种方式广播的大型或不可序列化变量(如KafkaProducer),但由于没有可见的RDD变量,我无法检索用于广播惰性求值变量的 Spark 上下文。如果我使用初始上下文来创建 DStreams 或从 DStreams 中检索到的上下文,则任务变得不可序列化。

那么如何在映射函数中使用广播变量呢?或者是否有任何解决方法可以在映射函数中使用大型或不可序列化的变量?

【问题讨论】:

    标签: java apache-kafka spark-streaming


    【解决方案1】:

    我终于找到了解决办法。要使用这些功能,请使用转换函数而不是映射函数。在 transform 函数中,我们手动处理 RDD 并对其应用 map 函数,因此我们可以获取 RDD 的引用,从而从中获取 Spark 上下文。

    【讨论】:

      猜你喜欢
      • 2016-08-18
      • 2015-08-26
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2015-04-18
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多