【问题标题】:Replaying an RDD in spark streaming to update an accumulator在火花流中重放 RDD 以更新累加器
【发布时间】:2015-12-14 17:47:51
【问题描述】:

我实际上已经没有选择了。 在我的火花流应用程序中。我想保持一些键的状态。我从卡夫卡那里得到事件。然后我从事件中提取密钥,比如用户 ID。当没有来自 Kafka 的事件时,我想每 3 秒更新一次相对于每个用户 ID 的计数器,因为我将 StreamingContext 的批处理持续时间配置为 3 秒。

现在我这样做的方式可能很难看,但至少它有效:我有一个这样的 accumulableCollection:

val userID = ssc.sparkContext.accumulableCollection(new mutable.HashMap[String,Long]())

然后我创建一个“假”事件并继续将其推送到我的 spark 流上下文,如下所示:

val rddQueue = new mutable.SynchronizedQueue[RDD[String]]()
for ( i <- 1 to  100) {
  rddQueue += ssc.sparkContext.makeRDD(Seq("FAKE_MESSAGE"))
  Thread.sleep(3000)
}
val inputStream = ssc.queueStream(rddQueue)

inputStream.foreachRDD( UPDATE_MY_ACCUMULATOR )

这将使我能够访问我的 accumulatorCollection 并更新所有用户 ID 的所有计数器。到目前为止,一切正常,但是当我从以下位置更改循环时:

for ( i <- 1 to  100) {} #This is for test

收件人:

while (true) {} #This is to let me access and update my accumulator through the whole application life cycle

然后当我运行 ./spark-submit 时,我的应用程序卡在这个阶段:

15/12/10 18:09:00 INFO BlockManagerMasterActor: Registering block manager slave1.cluster.example:38959 with 1060.3 MB RAM, BlockManagerId(1, slave1.cluster.example, 38959)

关于如何解决这个问题的任何线索?有没有一种非常简单的方法可以让我更新我的用户 ID 的值(而不是创建一个无用的 RDD 并定期将其推送到队列流)?

【问题讨论】:

    标签: apache-spark spark-streaming


    【解决方案1】:

    while (true) ... 版本不起作用的原因是控件永远不会返回到主执行行,因此该行以下的任何内容都不会被执行。为了解决这个特定问题,我们应该在单独的线程中执行while 循环。 Future { while () ...} 应该可以工作。 此外,在上面的示例中填充QueueDStream 时不需要Thread.sleep(3000)。 Spark Streaming 将在每个流式传输间隔消耗队列中的一条消息。

    触发“tick”消息流入的更好方法是使用ConstantInputDStream,它在每个流式传输间隔播放相同的RDD,因此无需使用QueueDStream 创建RDD 流入。

    也就是说,在我看来,目前的方法似乎很脆弱,需要修改。

    【讨论】:

      猜你喜欢
      • 2016-04-23
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2017-02-28
      • 1970-01-01
      • 2023-03-13
      • 1970-01-01
      相关资源
      最近更新 更多