【发布时间】:2019-01-06 13:16:31
【问题描述】:
我尝试在 Monix 中按键拆分单个 Observable,然后将每个 GrouppedObservable 中的最后一个 n 事件分组并发送它们以进行进一步处理。问题是要分组的键的数量可能是无限的,这会导致内存泄漏。
应用上下文:
我有来自许多对话的消息的 kafka 流。每个对话都有roomId,我想对这个ID进行分组以获取Observables的集合,每个对话只包含来自单个对话的消息。
会话室通常是短暂的,即用唯一的roomId创建新会话,在短时间内交换几十条消息,然后关闭会话。
为避免内存泄漏,我希望仅保留 100-1000 个最近对话的缓冲区,并删除较旧的对话。因此,如果一个事件来自一个长期未见的对话,它将被视为新对话,因为其先前消息的缓冲区将被遗忘。
Monix 中的groupBy 方法有参数keysBuffer,指定如何处理关键缓冲区。
我认为将keyBuffer 指定为DropOld 策略将使我能够实现我想要的行为。
以下是所描述用例的简化版本。
import monix.execution.Scheduler.Implicits.global
import monix.reactive._
import scala.concurrent.duration._
import scala.util.Random
case class Event(key: Key, value: String, seqNr: Int) {
override def toString: String = s"(k:$key;s:$seqNr)"
}
case class Key(conversationId: Int, messageNr: Int)
object Main {
def main(args: Array[String]): Unit = {
val fakeConsumer = Consumer.foreach(println)
val kafkaSimulator = Observable.interval(1.millisecond)
.map(n => generateHeavyEvent(n.toInt))
val groupedMessages = kafkaSimulator.groupBy(_.key)(OverflowStrategy.DropOld(50))
.mergeMap(slidingWindow)
groupedMessages.consumeWith(fakeConsumer).runSyncUnsafe()
}
def slidingWindow[T](source: Observable[T]): Observable[Seq[T]] =
source.scan(List.empty[T])(fixedSizeList)
def fixedSizeList[T](list: List[T], elem: T): List[T] =
(list :+ elem).takeRight(5)
def generateHeavyEvent(n: Int): Event = {
val conversationId: Int = n / 500
val messageNr: Int = n % 5
val key = Key(conversationId, messageNr)
val value = (1 to 1000).map(_ => Random.nextPrintableChar()).toString()
Event(key, value, n)
}
}
但是,在 VisualVM 上观察应用程序堆表明内存泄漏。跑了大约30分钟后,我得到了java.lang.OutOfMemoryError: GC overhead limit exceeded
下面是堆使用图的屏幕截图,描绘了我的应用运行了大约 30 分钟。 (最后扁平部分在OutOfMemoryError之后)
VisualVM Heap plot of application
我的问题是:如何在 monix 中通过可能无限数量的键对事件进行分组而不会泄漏内存?允许丢弃旧密钥
背景信息:
- monix 版本:
3.0.0-RC2 - scala 版本:
2.12.8
【问题讨论】:
-
你不想用
groupBy(_.key.conversationId)代替_.key吗?