【问题标题】:How to clean up substreams in continuous Akka streams如何清理连续 Akka 流中的子流
【发布时间】:2017-05-17 20:20:45
【问题描述】:

鉴于我有一个非常长的事件流,如下所示。当很长一段时间过去时,将创建许多不再需要的子流。

有没有办法在给定的时间清理特定的子流,因为 例如 id 3 创建的子流应该被清理并且状态 在 13Pm 丢失的扫描方法中(过期的 Wid 属性)?

case class Wid(id: Int, v: String, expires: LocalDateTime)
test("Substream with scan") {
  val (pub, sub) = TestSource.probe[Wid]
    .groupBy(Int.MaxValue, _.id)
    .scan("")((a: String, b: Wid) => a + b.v)
    .mergeSubstreams
    .toMat(TestSink.probe[String])(Keep.both)
    .run()
}

【问题讨论】:

    标签: scala akka akka-stream


    【解决方案1】:

    TL;DR 您可以在一段时间后关闭子流。但是,使用输入通过内置阶段动态设置时间是另一回事。

    关闭子流

    要关闭流程,您通常会(从上游)完成它,但您也可以(从下游)取消它。例如,take(n: Int) 流将在 n 元素通过后取消。

    现在,在groupBy 的情况下,您无法完成子流,因为上游由所有子流共享,但您可以取消它。如何取决于您要为其添加什么条件。

    但是,请注意groupBy 会删除已关闭的子流的输入:如果在 3-substream 关​​闭后,具有 id 3 的新元素从上游来到 groupBy,它将被简单地忽略并拉入下一个元素。其原因可能是在关闭和重新打开子流之间的过程中可能会丢失某些元素。此外,如果您的流应该运行很长时间,这将影响性能,因为在转发到相关(实时)子流之前,每个元素都将根据关闭的子流列表进行检查。如果您对它的性能不满意,您可能想要实现自己的有状态过滤器(例如,使用布隆过滤器)。

    要关闭子流,我通常使用take(如果您只想要给定数量的元素,但在无限流上可能不是这种情况),或者某种超时:completionTimeout 如果您想要从实现到关闭的固定时间或idleTimeout,如果您想在一段时间内没有元素通过时关闭。请注意,这些流不会取消流而是使其失败,因此您必须使用recoverrecoverWith 阶段捕获异常以将失败更改为取消(recoverWith 允许您在不发送任何最后一个元素的情况下取消,通过Source.empty 恢复)。

    动态设置超时时间

    现在你想要的是根据第一个通过的元素动态设置关闭时间。这更复杂,因为流的物化独立于通过它们的元素。事实上,在通常的情况下(没有groupBy),流在任何元素通过它们之前就被物化了,因此使用元素来物化它们是没有意义的。

    我在that question 中遇到了类似的问题,最终使用了带有签名的groupBy 的修改版本

    paramGroupBy[K, OO, MM](maxSubstreams: Int, f: Out => K, paramSubflow: K => Flow[Out, OO, MM])
    

    允许使用定义它的键定义每个子流。这可以修改为将第一个元素(而不是键)作为参数。

    另一种(在您的情况下可能更简单)方法是编写您自己的阶段,该阶段完全符合您的要求:从第一个元素获取结束时间并在那时取消流。这是一个示例实现(我使用调度程序而不是设置状态):

    object CancelAfterTimer
    
    class CancelAfter[T](getTimeout: T => FiniteDuration) extends GraphStage[FlowShape[T, T]] {
      val in = Inlet[T]("CancelAfter.in")
      val out = Outlet[T]("CancelAfter.in")
      override val shape: FlowShape[T, T] = FlowShape(in, out)
    
      override def createLogic(inheritedAttributes: Attributes): GraphStageLogic = new TimerGraphStageLogic(shape) with InHandler with OutHandler  {
        override def onPush(): Unit = {
          val elem = grab(in)
          if (!isTimerActive(CancelAfterTimer))
            scheduleOnce(CancelAfterTimer, getTimeout(elem))
          push(out, elem)
        }
    
        override def onTimer(timerKey: Any): Unit = 
          completeStage() //this will cancel the upstream and close the downstrean
    
        override def onPull(): Unit = pull(in)
    
        setHandlers(in, out, this)
      }
    }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2021-12-08
      • 2020-03-22
      • 2011-09-27
      • 2017-07-17
      • 2016-01-12
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多