【发布时间】: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