【发布时间】:2016-12-24 21:03:49
【问题描述】:
我打算基于一些网络数据的会话来实现一个火花流应用程序。我正在为 RDD 使用全状态编程。
由于大量的记录和键,我需要在我的流逻辑中的某些条件被遵守后删除我的 mapwithsate 函数中的一些状态!
我想知道为什么要这样做,我知道在 sate 规范中,有一个超时,但这不是我正在寻找的功能,而是我应该从内存中删除状态以减轻我的流式传输的内存量应用程序消耗。
例如下面是一个示例饱和函数
def trackStateFunc(batchTime: Time, key: String, value: Option[Int],state: State[Long]): Option[(String, Long)] = {
val sum = value.getOrElse(0).toLong + state.getOption.getOrElse(0L)
val output = (key, sum)
state.update(sum)
Some(output)
}
如果应用了某些条件,我想知道如何删除键的状态,以便释放我的流式应用程序所需的内存..
【问题讨论】:
标签: apache-spark spark-streaming stateful