拓冰建站拓冰建站
首页 / 资讯中心 / 正文

Akka Streams 的 Source.never:永不发射、永不完成、永不失败的无限等待数据源

后端并发编程异步编程【免费下载链接】akka-coreA platform to build and run apps that are elastic, agile, and resilient. SDK, libraries, and hosted environments.项目地址https://gitcode.com/gh_mirrors/ak/akka-core点击查看免费下载导读Source.never是 Akka Streams 中一个极为特殊的数据源Source它不发射任何元素、永不完成complete也永不失败fail。本文基于 Akka 官方运算符文档never.md深入讲解该运算符的 API 签名、Reactive Streams 语义、底层 GraphStage 实现并结合仓库中的源码与测试用例说明其典型应用场景——尤其是测试中模拟无限等待的下游行为以及它与Source.empty、Sink.never之间的区别与配合。一、Source.never 是什么从官方文档的定义来看Never emit any elements, never complete and never fail.Source.never创建的数据源具有三个确定性的行为特征不发射emits任何元素——下游无论请求多少次都拿不到数据永不完成completes——流不会正常结束onComplete永远不会触发永不失败fails——流也不会以异常方式终止onError同样永远不会触发。官方文档明确指出其用途Useful for tests。它常被用来模拟一个始终不给出结果的上游从而验证下游在超时、空闲或取消cancel场景下的行为是否符合预期。二、API 签名Scala DSLSource.never定义在 Source.scaladef never[T]: Source[T, NotUsed] _never private[this] val _never: Source[Nothing, NotUsed] fromGraph(GraphStages.NeverSource)签名说明泛型参数T由调用点推断因此Source.never[Int]、Source.never[String]均可使用返回类型为Source[T, akka.NotUsed]即物化值materialized value为NotUsed——该数据源在运行时不产生任何有意义的运行结果NotUsed仅表示无值底层通过fromGraph(GraphStages.NeverSource)构造直接复用同一个单例图private[this] val _never说明这是一个零状态、可安全复用的纯惰性数据源。调用示例import akka.NotUsed import akka.actor.ActorSystem import akka.stream.scaladsl.Source implicit val system: ActorSystem ActorSystem(never-example) val neverSource: Source[Int, NotUsed] Source.never[Int]Java DSLJava API 同样提供对应方法定义在 javadsl/Source.scalaimport akka.NotUsed; import akka.stream.javadsl.Source; SourceInteger, NotUsed neverSource Source.never();Java 版实现只是对 Scala 版本的简单包装scaladsl.Source.never.asJava。官方 apidoc 中 Java 签名为never()泛型由变量声明处的类型推断。三、Reactive Streams 语义官方文档以 callout 形式给出了规范化的语义描述行为结果emits发射never永不发射任何元素completes完成never永不完成这意味着当下游对never数据源发出request(n)需求信号后上游永远不回应任何元素而当流被物化运行后它也永远不触发完成或失败回调。整个流会保持挂起状态直到被下游主动cancel()或整个 ActorSystem 被终止。这种语义可以对照 Source.emptyempty是立即完成且不发射任何元素二者都不发射元素区别在于empty立即完成而never永不完成。二者互为补充分别覆盖需要空流与需要无限等待流两种测试场景。四、底层实现GraphStages.NeverSourceSource.never的实现核心是akka.stream.impl.fusing.GraphStages中的NeverSource单例见 GraphStages.scalaprivate[akka] object NeverSource extends GraphStage[SourceShape[Nothing]] { private val out OutletNothing val shape: SourceShape[Nothing] SourceShape(out) override def initialAttributes: Attributes DefaultAttributes.neverSource override def createLogic(inheritedAttributes: Attributes): GraphStageLogic with OutHandler new GraphStageLogic(shape) with OutHandler { override def onPull(): Unit () setHandler(out, this) } }从源码结构可以清晰看出永不发射的实现原理它是一个GraphStage[SourceShape[Nothing]]输出口类型为NothingNothing 是所有类型的子类型因此可被推断为任意T这解释了为什么Source.never[T]可以适用于任意元素类型它只实现了OutHandler.onPull()且方法体为空()——当下游请求元素时NeverSource什么都不做既不推元素、也不完成、也不报错它没有实现onDownstreamFinish之外的任何完成/失败逻辑因此从语义上讲就是永远等待。其initialAttributes指向DefaultAttributes.neverSource见 Stages.scala 处的val neverSource name(neverSource)这让调试时可以在运算符图谱中识别出该阶段。五、配套运算符Sink.never与Source.never配套Akka Streams 还提供了 Sink.neverdef never: Sink[Any, Future[Done]] _never它的物化值是Future[Done]该Future在上游完成时成功、在上游失败时失败而一旦有元素被推入NeverSink会直接以IllegalStateException(NeverSink should not receive any push.)失败见 GraphStages.scala 中的NeverSink实现。在测试中Source.never与Sink.never常被组合使用来验证流保持运行但不产生任何数据的场景。六、测试用例验证仓库中专门为Source.never编写了测试 NeverSourceSpec.scala完整验证了其核心语义The Never Source must { never completes in { val neverSource Source.never[Int] val pubSink Sink.asPublisherInt val neverPub neverSource.toMat(pubSink)(Keep.right).run() val c TestSubscriber.manualProbe[Int]() neverPub.subscribe(c) val subs c.expectSubscription() subs.request(1) c.expectNoMessage(300.millis) // 请求 1 个元素但 300ms 内没有任何消息 subs.cancel() } }该测试的关键步骤将Source.never[Int]接到一个Sink.asPublisher上并运行得到上游的 Publisher用TestSubscriber.manualProbe手动订阅并request(1)请求一个元素断言expectNoMessage(300.millis)——即使下游发出了需求信号300ms 内也观察不到任何元素、完成或失败信号这正是never emits / never completes / never fails语义的实证最后cancel()主动取消订阅避免测试挂起。这也提示了使用Source.never的注意事项它不会自行终止在真实业务代码中若不加超时控制或取消逻辑流将无限期等待因此它几乎专用于测试环境。七、典型应用场景与实践建议综合官方文档与源码Source.never的典型应用场景包括超时/空闲行为测试作为上游接入下游验证completionTimeout、idleTimeout等超时类运算符在长时间无数据时是否按预期触发或验证merge、concat、interleave等合并运算符在某个输入永不产出时的表现取消传播测试验证下游cancel()信号能否正确沿流向上游传播并释放资源占位数据源在需要提供一个 Source 但暂时不准备产出任何数据的 API 场景中充当占位实现配合Source.empty区分立即结束与永久等待两种语义。实践建议物化值约定Source.never的物化值是NotUsed如需获得可观测的运行句柄应搭配其他可物化出有用值的 Sink如Sink.asPublisher、Sink.actorRef必须主动取消由于流永不完成、永不失败测试结束时务必cancel()或依赖测试框架的 ActorSystem 清理机制避免资源泄漏Java/Scala 通用Java 用户直接调用Source.never()即可获得等价行为API 语义与 Reactive Streams 语义在两种 DSL 中完全一致。相关文档运算符总览stream/operators/index.md对照运算符Source.empty立即完成、不发射元素的空源核心实现GraphStages.NeverSourceScala APISource.neverJava APIjavadsl Source.never测试用例NeverSourceSpec.scala赞分享后端并发编程异步编程【免费下载链接】akka-coreA platform to build and run apps that are elastic, agile, and resilient. SDK, libraries, and hosted environments.项目地址https://gitcode.com/gh_mirrors/ak/akka-core点击查看免费下载相关推荐Akka Streams 的 Sink.never 详解永不消费、永不取消的背压型 SinkAkka Streams 的 Sink.never 详解永不消费、永不取消的背压型 Sink Sink.never 是 Akka Streams 提供的一个特后端并发编程异步编程Akka Streams Source.empty 详解立即完成且不发射任何元素的空数据源Akka Streams Source.empty 详解立即完成且不发射任何元素的空数据源 导读 Source.empty 是 Akka Streams 中最后端并发编程异步编程OpenJK性能优化揭秘为什么你的绝地学院运行更流畅了OpenJK性能优化揭秘为什么你的绝地学院运行更流畅了 OpenJK作为《星球大战绝地学院》和《绝地放逐者》的社区维护项目通过一系列深度优化让这款经典游戏游戏开发创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门