对Akka Stream中Source.delay方法及KillSwitch文档示例存疑且运行报错
关于Akka Stream KillSwitch.shutdown示例结果的疑惑解答
我来帮你拆解这个问题,先理清楚示例的预期逻辑,再解释你运行时遇到异常的原因:
为什么示例预期结果是2?
这个示例的核心是展示KillSwitch.shutdown()的优雅关闭特性:它不会立刻截断流,而是允许已经进入处理管道的元素完成传递后再终止流。示例的预期时序大致是这样的:
- 流启动后,
Sink.last会向上游发送元素需求,上游立刻产生元素1并进入delay阶段,开始1秒的计时。 - 假设
doSomethingElse()方法刚好执行了约1秒,此时元素1的延迟完成,被发送到Sink.last;同时,因为Sink.last会持续请求后续元素(它需要捕获最后一个元素),上游会继续产生元素2并进入delay阶段,开始新的1秒计时。 - 此时调用
killSwitch.shutdown():这个信号会通知上游停止产生新元素,但已经进入delay阶段的元素2会继续完成延迟计时,之后被发送到Sink.last。 - 流正常终止后,
Sink.last返回最后收到的元素2,所以预期结果是2。
你运行时出现异常的原因
你遇到的NoSuchElementException: last of empty stream,是因为doSomethingElse()的执行时间太短——在第一个元素1的1秒延迟完成前,你就调用了killSwitch.shutdown()。此时的时序是:
- 元素
1进入delay阶段开始计时,但还没到1秒的延迟时间。 killSwitch.shutdown()触发流的终止信号,对于delay操作来说,当收到终止信号时,还在等待延迟的元素会被直接丢弃(这是DelayOverflowStrategy.backpressure下的终止行为)。Sink.last从未收到任何元素,当流终止时自然会抛出空流异常。
如果想复现示例的预期结果,你可以在doSomethingElse()中加入1秒左右的延迟,比如:
def doSomethingElse(): Unit = Thread.sleep(1000)
这样就能让第一个元素完成延迟并传递到Sink,同时第二个元素进入delay阶段,此时shutdown会让第二个元素也完成传递。
内容的提问来源于stack exchange,提问作者Alexander Arendar
相关产品推荐
相关产品推荐

