TL;DR 您可以在一段时间后关闭子流。但是,使用输入通过内置阶段动态设置时间是另一回事。
关闭子流
要关闭流程,您通常会(从上游)完成它,但您也可以(从下游)取消它。例如,take(n: Int) 流将在 n 元素通过后取消。
现在,在groupBy 的情况下,您无法完成子流,因为上游由所有子流共享,但您可以取消它。如何取决于您要为其添加什么条件。
但是,请注意groupBy 会删除已关闭的子流的输入:如果在 3-substream 关闭后,具有 id 3 的新元素从上游来到 groupBy,它将被简单地忽略并拉入下一个元素。其原因可能是在关闭和重新打开子流之间的过程中可能会丢失某些元素。此外,如果您的流应该运行很长时间,这将影响性能,因为在转发到相关(实时)子流之前,每个元素都将根据关闭的子流列表进行检查。如果您对它的性能不满意,您可能想要实现自己的有状态过滤器(例如,使用布隆过滤器)。
要关闭子流,我通常使用take(如果您只想要给定数量的元素,但在无限流上可能不是这种情况),或者某种超时:completionTimeout 如果您想要从实现到关闭的固定时间或idleTimeout,如果您想在一段时间内没有元素通过时关闭。请注意,这些流不会取消流而是使其失败,因此您必须使用recover 或recoverWith 阶段捕获异常以将失败更改为取消(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)
}
}