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),仅供参考

相关新闻

ISO 9001 2026版不存在!当前应聚焦价值流与韧性管理

ISO 9001 2026版不存在!当前应聚焦价值流与韧性管理

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/24 5:58:09 阅读更多 →
RRU级联组网全解析:室分组网优缺点与典型应用场景

RRU级联组网全解析:室分组网优缺点与典型应用场景

摘要:RRU级联是室分组网中节约光纤的主流方案,通过光纤将多个RRU串联接入BBU一个光口。其核心优势在于节省光纤资源、降低组网成本,但链路前端故障会导致后端RRU全部中断。本文对比级联、星型与环型三种组网方式,并梳理室内分布、…

2026/9/24 5:58:09 阅读更多 →
Linux不重启刷新分区表:partprobe、partx与设备rescan实战

Linux不重启刷新分区表:partprobe、partx与设备rescan实战

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/24 5:58:09 阅读更多 →

最新新闻

旧华为手机救砖降级实战:MRT HW Tool一键脚本避坑指南

旧华为手机救砖降级实战:MRT HW Tool一键脚本避坑指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/24 8:02:22 阅读更多 →
板载声卡跑ASIO实战:官方驱动实现低延迟直播唱歌

板载声卡跑ASIO实战:官方驱动实现低延迟直播唱歌

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/24 8:02:22 阅读更多 →
Kornia 2D 边界框形状校验修复:`infer_bbox_shape` 与 `bbox_to_mask` 拒绝 rank-4 批处理输入并抛出 `ShapeError`

Kornia 2D 边界框形状校验修复:`infer_bbox_shape` 与 `bbox_to_mask` 拒绝 rank-4 批处理输入并抛出 `ShapeError`

计算机视觉人工智能深度学习图像处理 【免费下载链接】kornia 🐍 Geometric Computer Vision Library for Spatial AI 项目地址: https://gitcode.com/gh_mirrors/ko/kornia 点击查看 免费下载 本文解读 Kornia 几何模块中的一项行为修复(对…

2026/9/24 8:02:22 阅读更多 →
WinPE运维实战:从删除顽固文件到离线杀毒的完整指南

WinPE运维实战:从删除顽固文件到离线杀毒的完整指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/24 8:02:22 阅读更多 →
轨道交通EMC屏蔽关键:铍铜弹片选型与安装实战指南

轨道交通EMC屏蔽关键:铍铜弹片选型与安装实战指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/24 8:02:22 阅读更多 →
LeetCode-1. 两数之和

LeetCode-1. 两数之和

这里写目录标题方法一 暴力法方法二Python3 字典 (dict) 学习笔记一、字典语法格式二、创建字典1. 创建空字典2. 普通字典创建三、访问字典的值1. [键]方式取值2. 安全取值 get ()四、修改字典update () 批量更新五、删除字典元素pop / popitem方法一 暴力法 class Solution:d…

2026/9/24 8:01:22 阅读更多 →

日新闻

基于YOLOv8的渔船作业监控系统:从环境搭建到边缘部署全流程

基于YOLOv8的渔船作业监控系统:从环境搭建到边缘部署全流程

简介:这是一套面向计算机、人工智能、自动化等专业学生与教师的毕业设计级项目资源,围绕YOLOv8实现渔船作业监控系统,可用于毕设、课程设计、大作业或项目立项演示。压缩包共97个文件,约24.21MB,以70个Python源码文件为…

2026/9/24 0:00:19 阅读更多 →
单细胞注释实战:基于Scanpy的标记基因与参考映射流程解析

单细胞注释实战:基于Scanpy的标记基因与参考映射流程解析

简介:一份基于单细胞RNA测序数据的细胞类型注释算法研究Python毕业设计源码,针对计算机相关专业正在做毕设或需要项目实战的学习者,可用于课程设计与期末大作业。项目代码完整、经导师指导评审通过,可直接运行,覆盖数据…

2026/9/24 0:00:19 阅读更多 →
C#源生成器实战:用增量生成器替代反射,告别AOT崩溃

C#源生成器实战:用增量生成器替代反射,告别AOT崩溃

第一次在项目里被反射卡住,是在一个老旧的WinForms模块里:几十个类依赖PropertyChanged通知,运行时反射读属性、发通知,每次启动慢半拍不说,一上.NET Native/AOT裁剪模式几乎全面崩盘。后来我把这段逻辑全部改成C#源生…

2026/9/24 0:00:19 阅读更多 →

周新闻

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

直接铺开项目本身吧。这几个月我一直在折腾一件事:用Flutter给OpenHarmony做一款游戏集合类的App,说白了就是把若干小游戏塞进一个壳里,用统一入口分发。这个方向本身不算新鲜,真正让我花了不少心思的,是首页那堆游戏卡…

2026/9/23 4:55:02 阅读更多 →
Word表格编号全攻略:从列表编号到题注交叉引用

Word表格编号全攻略:从列表编号到题注交叉引用

写Word文档,最让人头疼的往往是那些“看起来不起眼”的小问题。比如表格编号这事:今天在表后面多加了两个空白行,明天给客户交稿前发现整个章节的编号全部错位,光是挨个改序号就能耗掉大半个下午。我前阵子帮人整理一份上百页的技…

2026/9/23 4:49:06 阅读更多 →
从第一个站到第二个站:独立开发者的静态网站选型与落地实践

从第一个站到第二个站:独立开发者的静态网站选型与落地实践

1. 项目概述1.1 核心需求解析做独立开发者这几年,说实话,第一个网站上线的那天晚上我兴奋得没睡着。但等它跑了半年,流量惨淡、功能臃肿、代码自己都懒得看第二遍之后,我才慢慢琢磨明白一个道理:第一个网站是练手&…

2026/9/23 9:53:41 阅读更多 →

月新闻

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能分类:[AI/大模型]细分主题:AI 增强型 CI/CD 流水线自动化与 GitOps 实践:Agent 工作流、工具调用与任务拆解:从原型到生产的验收清单很多团队在尝试用大…

2026/9/23 9:53:40 阅读更多 →
容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场分类:[工程技术]细分主题:Kubernetes 生产环境运维与排障实战:可复制的项目复盘模板与决策记录大部分团队的事故复盘报告,最后都变成了躺在 Confluence 或钉…

2026/9/23 9:53:40 阅读更多 →
容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步分类:[工程技术]细分主题:Docker 容器化技术与镜像安全管理:核心链路的逐步实现与关键代码取舍面对一个积累了五六年历史包袱的单体架构应用(包含 Web 接口、后台…

2026/9/23 9:53:40 阅读更多 →