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