use*_*545 7 scala akka akka-stream
鉴于我有很长时间的事件流经过如下所示的事物.当很长一段时间过去了,将会有很多不再需要的子流.
有没有办法在给定时间清理特定的子流,例如,应清除由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()
}
Run Code Online (Sandbox Code Playgroud)
TL; DR您可以在一段时间后关闭子流.但是,使用输入来动态设置内置阶段的时间是另一回事.
关闭子流
要关闭流程,您通常会(从上游)完成它,但您也可以取消它(从下游).例如,take(n: Int)一旦n元素通过,流程将取消.
现在,在这种groupBy情况下,您无法完成子流,因为上游由所有子流共享,但您可以取消它.如何取决于你想要放在什么条件.
但是,要知道,groupBy对于那些已经关闭的子流删除输入:如果ID为一个新的元素3来自于上游的groupBy后3-substream已经关闭,它只会被忽略,下一个元素将被拉到原因.因为这可能是在关闭和重新打开子流之间的过程中可能会丢失一些元素.此外,如果您的流应该运行很长时间,这将影响性能,因为在转发到相关(实时)子流之前,将针对关闭的子流列表检查每个元素.如果您对此的性能不满意,您可能希望实现自己的有状态过滤器(例如,使用布隆过滤器).
要关闭一个子流,我通常使用take(如果你只想要一定数量的元素,但在无限流中可能不是这种情况),或者某种超时:要么completionTimeout你想要从物化到关闭的固定时间,要么idleTimeout如果你想在没有元素通过一段时间后关闭.请注意,这些流不会取消流但会使其失败,因此您必须使用recover或recoverWith阶段捕获异常以将失败更改为取消(recoverWith允许您通过恢复而取消而不发送任何最后一个元素Source.empty).
动态设置超时
现在你想要的是根据第一个传递元素动态设置关闭时间.这更复杂,因为流的实现与通过它们的元素无关.实际上,在通常(没有groupBy)的情况下,流在任何元素通过之前都已实现,因此使用元素实现它们是没有意义的.
我在那个问题上遇到了类似的问题,最后使用了groupBy带签名的修改版本
paramGroupBy[K, OO, MM](maxSubstreams: Int, f: Out => K, paramSubflow: K => Flow[Out, OO, MM])
Run Code Online (Sandbox Code Playgroud)
允许使用定义它的键定义每个子流.这可以修改为具有第一个元素(而不是键)作为参数.
另一种(可能更简单,在你的情况下)方式是编写你自己的阶段,完全符合你的要求:从第一个元素获取结束时间并在那时取消流.这是一个示例实现(我使用调度程序而不是设置状态):
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)
}
}
Run Code Online (Sandbox Code Playgroud)