【问题标题】:Can't access broadcast variable in transformation无法在转换中访问广播变量
【发布时间】:2016-06-24 00:05:07
【问题描述】:

我在从转换函数内部访问变量时遇到问题。有人可以帮我吗? 这是我的相关类和函数。

@SerialVersionUID(889949215L)
object MyCache extends Serializable {
    @transient lazy val logger = Logger(getClass.getName)
    @volatile var cache: Broadcast[Map[UUID, Definition]] = null

    def getInstance(sparkContext: SparkContext) : Broadcast[Map[UUID, Definition]] = {
        if (cache == null) {
            synchronized {
                val map = sparkContext.cassandraTable("keyspace", "table")
                   .collect()
                   .map(m => m.getUUID("id") ->
                        Definition(m.getString("c1"), m.getString("c2"), m.getString("c3"),
                                m.getString("c4"))).toMap
                cache = sparkContext.broadcast(map)
            }
        }
        cache
    }
}

在不同的文件中:

object Processor extends Serializable {
    @transient lazy val logger = Logger(getClass.getName)

    def processData[T: ClassTag](rawStream: DStream[(String, String)], ssc: StreamingContext,
                                        processor: (String, Broadcast[Map[UUID, Definition]]) => T): DStream[T] = {
        MYCache.getInstance(ssc.sparkContext)
        var newCacheValues = Map[UUID, Definition]()
        rawStream.cache()
        rawStream
          .transform(rdd => {
                val array = rdd.collect()
                array.foreach(r => {
                      val value = getNewCacheValue(r._2, rdd.context)
                      if (value.isDefined) {
                          newCacheValues = newCacheValues + value.get
                      }
                })
                rdd
           })
       if (newCacheValues.nonEmpty) {
           logger.info(s"Rebroadcasting.  There are ${newCacheValues.size} new values")
           logger.info("Destroying old cache")
           MyCache.cache.destroy()
           // this is probably wrong here, destroying object, but then referencing it.  But I haven't gotten to this part yet.
           MyCache.cache = ssc.sparkContext.broadcast(MyCache.cache.value ++ newCacheValues)
       }
       rawStream
          .map(r => {
               println("######################")
               println(MyCache.cache.value)
               r
          })
          .map(r => processor(r._2, MyCache.cache.value))
          .filter(r => null != r)
   }
}

每次运行此程序时,我都会在尝试访问 cache.value 时得到 SparkException: Failed to get broadcast_1_piece0 of broadcast_1

当我在 .getInstance 之后添加 println(MyCache.cache.values) 时,我可以访问广播变量,但是当我将它部署到 mesos 集群时,我无法再次访问广播值,但返回值为 null指针异常。

更新:

我看到的错误是在println(MyCache.cache.value)。我不应该添加这个包含破坏的 if 语句,因为我的测试从来没有达到这个目标。

我的应用程序的基础是,我在 cassandra 中有一个不会更新太多的表。但我需要对一些流数据做一些验证。所以我想把这张表中所有更新不多的数据都拉到内存中。 getInstance 在启动时将整个表拉入,然后我检查所有流数据以查看是否需要再次从 cassandra 中提取(我将很少这样做)。转换和收集是我检查是否需要拉入新数据的地方。但由于我的表可能会更新,所以我需要不时更新广播。所以我的想法是销毁它然后重播。一旦我让其他东西工作,我会更新它。

如果我注释掉销毁并重新广播,我会得到同样的错误。

另一个更新:

我需要在processor这一行访问广播变量:.map(r => processor(r._2, MyCache.cache.value))

我能够在转换中广播变量,如果我在转换中执行println(MyCache.cache.value),那么我的所有测试都会通过,然后我就可以访问processor 中的广播

更新:

rawStream
    .map(r => {
      println("$$$$$$$$$$$$$$$$$$$")
      println(metrics.value)
      r
    })

这是我到达这一行时得到的堆栈跟踪。

    ERROR org.apache.spark.executor.Executor - Exception in task 0.0 in stage 135.0 (TID 114)
    java.io.IOException: org.apache.spark.SparkException: Failed to get broadcast_1_piece0 of broadcast_1
        at org.apache.spark.util.Utils$.tryOrIOException(Utils.scala:1222)
        at org.apache.spark.broadcast.TorrentBroadcast.readBroadcastBlock(TorrentBroadcast.scala:165)
        at org.apache.spark.broadcast.TorrentBroadcast._value$lzycompute(TorrentBroadcast.scala:64)
        at org.apache.spark.broadcast.TorrentBroadcast._value(TorrentBroadcast.scala:64)
        at org.apache.spark.broadcast.TorrentBroadcast.getValue(TorrentBroadcast.scala:88)
        at org.apache.spark.broadcast.Broadcast.value(Broadcast.scala:70)
        at com.uptake.readings.ingestion.StreamProcessors$$anonfun$processIncomingKafkaData$4.apply(StreamProcessors.scala:160)
        at com.uptake.readings.ingestion.StreamProcessors$$anonfun$processIncomingKafkaData$4.apply(StreamProcessors.scala:158)
        at scala.collection.Iterator$$anon$11.next(Iterator.scala:370)
        at scala.collection.Iterator$$anon$11.next(Iterator.scala:370)
        at scala.collection.Iterator$$anon$13.hasNext(Iterator.scala:414)
        at org.apache.spark.storage.MemoryStore.unrollSafely(MemoryStore.scala:284)
        at org.apache.spark.CacheManager.putInBlockManager(CacheManager.scala:171)
        at org.apache.spark.CacheManager.getOrCompute(CacheManager.scala:78)
        at org.apache.spark.rdd.RDD.iterator(RDD.scala:268)
        at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:38)
        at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:306)
        at org.apache.spark.rdd.RDD.iterator(RDD.scala:270)
        at org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:73)
        at org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:41)
        at org.apache.spark.scheduler.Task.run(Task.scala:89)
        at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:214)
        at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142)
        at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617)
        at java.lang.Thread.run(Thread.java:745)
    Caused by: org.apache.spark.SparkException: Failed to get broadcast_1_piece0 of broadcast_1
        at org.apache.spark.broadcast.TorrentBroadcast$$anonfun$org$apache$spark$broadcast$TorrentBroadcast$$readBlocks$1$$anonfun$2.apply(TorrentBroadcast.scala:138)
        at org.apache.spark.broadcast.TorrentBroadcast$$anonfun$org$apache$spark$broadcast$TorrentBroadcast$$readBlocks$1$$anonfun$2.apply(TorrentBroadcast.scala:138)
        at scala.Option.getOrElse(Option.scala:121)
        at org.apache.spark.broadcast.TorrentBroadcast$$anonfun$org$apache$spark$broadcast$TorrentBroadcast$$readBlocks$1.apply$mcVI$sp(TorrentBroadcast.scala:137)
        at org.apache.spark.broadcast.TorrentBroadcast$$anonfun$org$apache$spark$broadcast$TorrentBroadcast$$readBlocks$1.apply(TorrentBroadcast.scala:120)
        at org.apache.spark.broadcast.TorrentBroadcast$$anonfun$org$apache$spark$broadcast$TorrentBroadcast$$readBlocks$1.apply(TorrentBroadcast.scala:120)
        at scala.collection.immutable.List.foreach(List.scala:381)
        at org.apache.spark.broadcast.TorrentBroadcast.org$apache$spark$broadcast$TorrentBroadcast$$readBlocks(TorrentBroadcast.scala:120)
        at org.apache.spark.broadcast.TorrentBroadcast$$anonfun$readBroadcastBlock$1.apply(TorrentBroadcast.scala:175)
        at org.apache.spark.util.Utils$.tryOrIOException(Utils.scala:1219)
        ... 24 more

【问题讨论】:

  • 错误在哪一行?看起来你甚至评论过“这可能是错误的”,因为它是MyCache.cache.value 的第一次访问,它不应该工作。在transform 中调用rdd.collect() 对我来说也很奇怪。
  • 在我看来,您在这里走错了路。广播旨在将静态的东西(如不可变的地图)分发给所有工作人员以便快速访问。从外观上看,您正在尝试构建地图,您不应该为此使用广播。我同意@Alexey Romanov 的观点,调用rdd.collect 似乎很奇怪,因为整个rdd 然后在驱动程序上处理而不是并行处理,并行处理是Spark 擅长的......
  • 糟糕,抱歉,我将更新我的问题。我在这里没有给出太多的上下文。

标签: scala apache-spark broadcast


【解决方案1】:

[更新答案]

您收到错误是因为rawStream.map 中的代码,即MyCache.cache.value 正在其中一个执行程序上执行,而MyCache.cache 仍然是null

当您执行MyCache.getInstance 时,它会在驱动程序上创建MyCache.cache 并广播它。但是您在 map 方法中引用的不是同一个对象,因此它不会被发送给执行者。相反,由于您直接引用 MyCache,因此执行程序在自己的 MyCache 对象副本上调用 MyCache.cache,这显然是 null。

您可以通过首先在驱动程序中获取cache 广播对象的实例并在地图中使用that 对象来使其按预期工作。以下代码应该适合您--

val cache = MYCache.getInstance(ssc.sparkContext)
rawStream.map(r => {
                     println(cache.value)
                     r
             })

【讨论】:

  • 糟糕,抱歉,我没有提供太多背景信息,让我更新一下我的问题。
  • 我想我现在可以看到这个问题了。 (我是 Stackoverflow 的新手。不确定更新答案的可接受程序是什么——但我会更新原始答案。)
  • 非常感谢。这很有意义。但是,我怎样才能在其他文件中引用它呢?我认为将它作为属性添加到这个静态类上会做同样的事情吗?这样我就不必将它传递给我想要引用它的每个函数。
  • 在驱动程序本身中,您仍然可以始终从 any-where/file 调用 val cache = MYCache.getInstance(ssc.sparkContext)。只要确保在执行器上执行某些代码时,您传递的是实际的广播对象,而不仅仅是静态方法调用。
  • 哦,这很有意义。我会尝试一下,如果可行,请回来接受这个答案。谢谢!
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2012-04-14
  • 2022-01-27
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多