Akka Streams滑动窗口实现问题:上游关闭Outlet不可用且测试超时
嘿,我来帮你梳理下你遇到的这个头疼问题:你实现了一个触发器驱动的滑动窗口,用Akka Streams TestKit写测试时,测试死活跑不完——onUpstreamFinish里的while循环卡着退不出来,最后还爆了超时的断言错误:
Exception in thread "main" java.lang.AssertionError: assertion failed: timeout (3 seconds) during expectMsg while waiting for OnNext(Stream(12, ?)) at scala.Predef$.assert(Predef.scala:219) at akka.testkit.TestKitBase.expectMsg_int...
我从几个常见的坑点给你排查方向和解决方案:
1. 先盯死onUpstreamFinish的循环终止条件
你说循环退不出来,十有八九是这个循环的停止逻辑有问题。比如:
- 缓冲区里还有剩余元素,但你没在流结束时主动触发最后一次窗口输出,导致循环一直空等
- 循环的判断条件依赖某个状态,但这个状态在upstream finish时根本没被正确更新
举个典型的错误写法,你踩的可能就是这个坑:
override def onUpstreamFinish(): Unit = { // 坏例子:一直检查缓冲区非空,但完全没处理/输出元素,死循环了 while (buffer.nonEmpty) { // 啥也不干,就等着缓冲区变空?不可能啊! } completeStage() }
正确的姿势是,在流结束时主动处理缓冲区里剩余的有效窗口内容,然后终止循环:
override def onUpstreamFinish(): Unit = { // 输出最后一批符合窗口规则的元素(比如滑动窗口的剩余片段) if (buffer.size >= windowSize) { push(Chunk(buffer.toList)) buffer.clear() } // 不管有没有剩余元素,最后一定要调用completeStage()结束阶段 completeStage() }
2. 测试代码没处理流的终止信号
Akka Streams TestKit的TestSink必须收到OnComplete信号才会认为流真的结束了。如果你的自定义窗口阶段没在upstream finish时调用completeStage(),测试就会一直傻等,直到超时。
检查下你的测试代码,是不是漏了这两步:
- 测试源发送完所有元素后调用了
complete() - 测试Sink最后调用了
expectComplete(),而不是只等OnNext
给你个正确的测试写法参考:
val testSource = Source(List(1,2,3,4,5)) val windowStage = new TriggeredSlidingWindowStage(windowSize = 2, slideStep = 1) testSource.via(windowStage) .runWith(TestSink.probe[List[Int]]) .request(100) // 先请求足够多的元素 .expectNext(List(1,2), List(2,3), List(3,4), List(4,5)) // 预期的窗口输出 .expectComplete() // 关键!要等流结束的信号
3. 触发器的线程安全问题要重视
既然是触发器驱动的窗口,得确保触发器的信号和流的处理逻辑是线程安全的。比如:
- 流已经结束了,但触发器还在不停发信号,导致缓冲区状态一直变,循环永远退不出来
- 没在
onUpstreamFinish里停止触发器的监听,后续信号一直在干扰
解决方法很简单,先停触发器,再处理缓冲区:
override def onUpstreamFinish(): Unit = { // 先把触发器的订阅停了,别让它再瞎发信号 trigger.unsubscribe(this) // 处理剩余的窗口元素 flushRemainingBuffer() // 最后结束阶段 completeStage() }
4. 临时调大超时(应急用,别依赖)
如果你的窗口处理逻辑确实需要更长时间,可以临时调大TestKit的超时时间救急:
// 全局设置超时 implicit val timeout: Timeout = Timeout(10.seconds) // 或者单独给expectMsg指定超时 probe.expectMsg(10.seconds, OnNext(Stream(12, ...)))
但这只是权宜之计,核心还是要把循环退不出的根因解决掉。
总结一下:你现在的核心问题就是onUpstreamFinish里的循环没有合理的终止条件,或者窗口阶段没在合适时机结束。先把循环里的剩余元素处理逻辑补全,确保循环能正常退出,再检查测试代码有没有等流结束的信号,应该就能解决了。
内容的提问来源于stack exchange,提问作者manu

