【发布时间】:2019-04-28 19:47:12
【问题描述】:
我正在学习 LMAX Disruptor 并遇到一个问题:当我有一个非常大的环形缓冲区(例如 1024)并且我的生产者比我的消费者快得多时,环形缓冲区将保存大量数据,但不会发布事件,直到我的应用程序结束。这意味着我的应用程序将丢失大量数据(我的应用程序不是守护程序)。
我试图降低生产者的速度,这很有效。但是我不能在我的应用程序中使用这种方法,它会大大降低我的应用程序的性能。
val ringBufferSize = 1024
val disruptor = new Disruptor[util.Map[String, Object]](new MessageEventFactory, ringBufferSize, new MessageThreadFactory, ProducerType.MULTI, new BlockingWaitStrategy)
disruptor.handleEventsWith(new MessageEventHandler(batchSize, this))
disruptor.setDefaultExceptionHandler(new MessageExceptionHandler)
val ringBuffer = disruptor.start
val producer = new MessageEventProducer(ringBuffer)
part.foreach { row =>
// Thread.sleep(2000)
accm.add(1)
producer.onData(row)
// flush(row)
}
我想找到一种方法来自己控制disruptor的batch size,有没有什么方法可以消耗我的应用程序结束时保存的其余数据?
【问题讨论】:
-
你见过this question吗? @jasonk 建议的方法之一可能会在您的情况下使用。 // 另外,如果您的主要问题是生产者在整个“批次”被消耗之前无法发布第 1025 个事件,您可能需要查看EarlyReleaseHandler 示例。
-
@Michael Barker 的This answer 还展示了一个较短的
SequenceReportingEventHandler实现示例。