检测Source.queue

如何解决检测Source.queue

我希望有一个Source.queue(或类似的东西可以将其推送到物化图上),用来告诉我队列的当前饱和度。

我想这样做,而不(重新)实现QueueSource图形阶段提供的功能

我想出的一种可能的解决方法是:

object InstrumentedSource {

  final class InstrumentedSourceQueueWithComplete[T](
      delegate: SourceQueueWithComplete[T],bufferSize: Int,)(implicit executionContext: ExecutionContext)
      extends SourceQueueWithComplete[T] {

    override def complete(): Unit = delegate.complete()

    override def fail(ex: Throwable): Unit = delegate.fail(ex)

    override def watchCompletion(): Future[Done] = delegate.watchCompletion()

    private val buffered = new AtomicLong(0)

    private[InstrumentedSource] def onDequeue(): Unit = {
      val _ = buffered.decrementAndGet()
    }

    object BufferSaturationRatioGauge extends RatioGauge {
      override def getRatio: RatioGauge.Ratio = RatioGauge.Ratio.of(buffered.get(),bufferSize)
    }

    lazy val bufferSaturationGauge: RatioGauge = BufferSaturationRatioGauge

    override def offer(elem: T): Future[QueueOfferResult] = {
      val result = delegate.offer(elem)
      result.foreach {
        case QueueOfferResult.Enqueued =>
          val _ = buffered.incrementAndGet()
        case _ => // do nothing
      }
      result
    }
  }

  def queue[T](bufferSize: Int,overflowStrategy: OverflowStrategy)(
      implicit executionContext: ExecutionContext,materializer: Materializer,): Source[T,InstrumentedSourceQueueWithComplete[T]] = {
    val (queue,source) = Source.queue[T](bufferSize,overflowStrategy).preMaterialize()
    val instrumentedQueue = new InstrumentedSourceQueueWithComplete[T](queue,bufferSize)
    source.mapMaterializedValue(_ => instrumentedQueue).map { item =>
      instrumentedQueue.onDequeue()
      item
    }
  }

}

这似乎主要是通过一些手动测试来实现的(除了buffered最终最终与队列中的实际项目数保持一致,在我的情况下应该是可以的),但是我想知道是否有解决方案可以更好地利用我可能错过的内置功能

版权声明:本文内容由互联网用户自发贡献,该文观点与技术仅代表作者本人。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如发现本站有涉嫌侵权/违法违规的内容, 请发送邮件至 dio@foxmail.com 举报,一经查实,本站将立刻删除。

相关推荐


Selenium Web驱动程序和Java。元素在(x,y)点处不可单击。其他元素将获得点击?
Python-如何使用点“。” 访问字典成员?
Java 字符串是不可变的。到底是什么意思?
Java中的“ final”关键字如何工作?(我仍然可以修改对象。)
“loop:”在Java代码中。这是什么,为什么要编译?
java.lang.ClassNotFoundException:sun.jdbc.odbc.JdbcOdbcDriver发生异常。为什么?
这是用Java进行XML解析的最佳库。
Java的PriorityQueue的内置迭代器不会以任何特定顺序遍历数据结构。为什么?
如何在Java中聆听按键时移动图像。
Java“Program to an interface”。这是什么意思?
Java在半透明框架/面板/组件上重新绘画。
Java“ Class.forName()”和“ Class.forName()。newInstance()”之间有什么区别?
在此环境中不提供编译器。也许是在JRE而不是JDK上运行?
Java用相同的方法在一个类中实现两个接口。哪种接口方法被覆盖?
Java 什么是Runtime.getRuntime()。totalMemory()和freeMemory()?
java.library.path中的java.lang.UnsatisfiedLinkError否*****。dll
JavaFX“位置是必需的。” 即使在同一包装中
Java 导入两个具有相同名称的类。怎么处理?
Java 是否应该在HttpServletResponse.getOutputStream()/。getWriter()上调用.close()?
Java RegEx元字符(。)和普通点?