首页/新闻资讯/正文详情

Akka Streams Source.lazily 算子详解:延迟 Source 创建与物化(含 lazySource 迁移指南)

发布时间:2026/9/24 15:53:11 来源:云帆数科 栏目:资讯中心
Akka Streams Source.lazily 算子详解:延迟 Source 创建与物化(含 lazySource 迁移指南)
后端并发编程异步编程【免费下载链接】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.lazily是 Akka Streams 中用于“延迟创建与物化 Source”的经典算子只有当下游产生需求demand时它才会调用工厂函数真正构造内部Source并完成物化。在 Akka 2.6.0 中该算子已被Source.lazySource取代本文以 lazily 算子文档 为主线结合 lazySource 算子文档、Scala DSL 源码、Java DSL 源码 与底层实现 LazySource.scala讲清其签名、语义、实现原理、典型用法与迁移路径帮助你正确使用这一延迟构造机制并规避“假懒加载”陷阱。一、核心概念把“创建”推迟到“有需求”那一刻lazily解决的是这样一个问题某些Source的构造过程非常昂贵例如建立数据库连接、打开文件、初始化重量级资源但并非每次流运行都必须真正用到它。如果直接构造该Source代价会在定义流图descriptor时或流物化时立即产生而用lazily包裹后工厂函数create的调用会被推迟到下游首次发出需求即onPull时。官方文档对它的语义描述只有一句话却是理解全部行为的关键Defers creation and materialization of aSourceuntil there is demand. 延迟Source的创建与物化直到产生需求为止。也就是说lazily不仅推迟了create工厂的执行还推迟了被包裹Source的物化materialization——两个步骤都会在第一次pull到达时一次性完成。关于该行为更完整的描述含失败与取消场景可参阅其继任者文档 lazySource.md。二、方法签名与弃用状态2.1 版本状态lazily在Akka 2.6.0起被标记为弃用deprecated官方明确要求改用Source.lazySource。这一弃用标记同时存在于 Scala 与 Java 两个 DSL 的源码中Scala DSLSource.scala 第 498 行 标注deprecated(Use Source.lazySource instead, 2.6.0)Java DSLjavadsl/Source.scala 第 272 行 标注同样的弃用信息。2.2 Scala 签名deprecated(Use Source.lazySource instead, 2.6.0) def lazilyT, M Source[T, M]): Source[T, Future[M]]create一个无参工厂函数返回真正要延迟构造的Source[T, M]返回值外层Source[T, Future[M]]——元素的类型为T而物化值是一个Future[M]该Future会在内部 Source 真正被物化时以内部 Source 的物化值M完成。2.3 Java 签名Deprecated public T, M SourceT, CompletionStageM lazily(CreatorSourceT, M create)Java 侧接受akka.japi.function.Creator返回的物化值是对应 ScalaFuture的 Java 视图CompletionStageM。其实现只是对 Scala DSL 的一层薄封装见 javadsl/Source.scala 第 273-274 行scaladsl.Source.lazilyT, M create.create().asScala).mapMaterializedValue(_.asJava).asJava2.4 姊妹算子lazilyAsync除了lazily同一批被弃用的还有lazilyAsyncScala 侧见 Source.scala 第 510 行Java 侧见 javadsl/Source.scala 第 284 行它接收一个返回Future[T]的工厂并在有需求时才调用deprecated(Use Source.lazyFuture instead, 2.6.0) def lazilyAsyncT Future[T]): Source[T, Future[NotUsed]]它的实现其实就是lazily(() fromFuture(create()))——把“延迟构造一个 future 值”转化为“延迟构造一个单元素 Source”。对应的新 API 是Source.lazyFuture。三、Reactive Streams 语义lazily本身是一个纯包装算子不改变内部 Source 的流语义。官方文档用 callout 形式明确给出属性语义emits发射取决于被包裹的Sourcedepends on the wrappedSourcecompletes完成取决于被包裹的Sourcedepends on the wrappedSource这意味着外层流何时发射元素、何时正常完成完全由内部Source决定lazily只负责“何时创建”与“何时物化”不增删任何元素也不改变完成与失败的时机。与之对应lazySource 文档 给出的 semantics 完全相同。四、底层实现原理一个 GraphStage 的故事lazily在 2.6.0 之前就由akka.stream.impl.LazySource这个内部GraphStage实现见 LazySource.scalalazily与lazySource共用同一套机制——事实上lazySource的源码就是fromGraph(new LazySource(create))Source.scala 第 592-593 行。理解它的createLogicAndMaterializedValue逻辑LazySource.scala 第 34-90 行就能彻底弄明白延迟语义的边界1. 物化值是一个 Promise。创建 stage 逻辑时会先构造一个Promise[M]最终暴露给用户的Future[M]/CompletionStage[M]就是它的 future。2. 工厂在onPull中才被调用。OutHandler.onPull是下游首次发出需求的通知。此时才执行sourceFactory()若工厂抛异常NonFatalmatPromise以该异常失败并重新抛出流失败若成功则创建一个SubSinkInlet把内部 Source 物化到其中subFusingMaterializer.materialize(...)并把内部 Source 的物化值写入matPromisematPromise.trySuccess(matVal)若内部 Source 物化失败则取消 subSink、使 stage 失败并使matPromise失败。3. 下游提前取消/失败会触发NeverMaterializedException。如果在下游第一次pull之前下游就取消或失败了则onDownstreamFinish会让matPromise以NeverMaterializedException失败LazySource.scala 第 38-41 行。这个异常类型定义在 NeverMaterializedException.scala它的存在与 lazySource 文档中的描述完全对应“如果工厂未被调用则物化值以NeverMaterializedException失败”。4. 中途取消会向上游传播。切换 handler 之后如果下游在内部 Source 运行中途取消onDownstreamFinish会调用subSink.cancel(cause)并结束 stageLazySource.scala 第 59-63 行。5. 异常中止兜底。postStop中如果matPromise尚未完成则以AbruptStageTerminationException使其失败LazySource.scala 第 84-86 行。从源码结构看LazySource继承GraphStageWithMaterializedValue[SourceShape[T], Future[M]]LazySource.scala 第 27-28 行这是 Akka Streams 中“自定义 stage 自定义物化值”的标准组合模式也是它能向外暴露Future[M]这一额外物化值的原因。五、典型用法一次物化一个可变对象延迟创建最有价值的应用场景是保证每次物化都能拿到一个全新的可变对象。文档 lazySource.md 给出的例子很能说明问题假设有一个类似迭代器的IteratorLikeThing它提供thereAreMore是否还有元素和extractNext取出下一个元素并前进两个方法。如果把同一个实例直接放进Source.unfold那么同一条流被多次物化多次run()时这个可变实例会被所有运行中的流共享——这是不安全的。而用lazily/lazySource包裹后工厂每次物化都会执行一次从而为每次运行创建一个独立实例。对应 Scala 测试示例位于 akka-docs/src/test/scala/docs/stream/operators/source/Lazy.scalaval stream Source .lazySource { () val iteratorLike new IteratorLikeThing Source.unfold(iteratorLike) { iteratorLike if (iteratorLike.thereAreMore) Some((iteratorLike, iteratorLike.extractNext)) else None } } .to(Sink.foreach(println)) // each of the three materializations will have their own instance of IteratorLikeThing stream.run() stream.run() stream.run()Java 侧等价的完整示例在 akka-docs/src/test/java/jdocs/stream/operators/source/Lazy.java核心结构一致工厂内new IteratorLikeThing()配合Source.unfold逐元素推进。提示如果你持有的是真正的java.util.Iterator官方建议优先使用Source.fromIterator如果资源是“打开-读取-关闭”形态也可以考虑Source.unfoldResource——二者往往比lazily更贴合场景。六、重要警告别被“假懒加载”骗了官方文档特别强调了一个极易踩坑的事实异步边界asynchronous boundaries和许多算子会做预取pre-fetching从而比预期更早地触发需求让工厂提前执行。也就是说lazily并不保证“直到你显式 pull 才创建”。lazySource 文档 给出了反例用Sink.queue消费流时你可能预期只有调用queue.pull()才会创建昂贵的 Source但实际上Sink.queue自带缓冲在物化时就会立即发起需求于是昂贵的 Source 很快就被创建了// #not-a-good-example —— 别以为这里真的懒 val source Source.lazySource { () println(Creating the actual source) createExpensiveSource() } val queue source.runWith(Sink.queue()) // ... time passes ... // at some point in time we pull the first time // but the source creation may already have been triggered queue.pull()因此如果绝对要在“第一次真正 pull”时才创建资源请确认整条链路中不存在任何预取/缓冲算子如果只是想避免在流定义阶段就付出创建成本以及让每次物化都得到独立的新实例那么即使工厂被提前触发lazily依然满足需求——这也是它最稳妥的使用方式。Scala 源码注释Source.scala 第 494-496 行与 lazySource 的源码注释Source.scala 第 569-571 行都明确记录了“异步边界与预取会抵消延迟效果”这一限制。七、从 lazily 迁移到 lazySource由于lazily在 2.6.0 起弃用新代码应直接使用lazySource老代码迁移成本极低——签名、语义、物化值类型完全一致只是改名弃用 API替代 API说明Source.lazily(create)Source.lazySource(create)延迟构造Source[T, M]物化值为Future[M]Source.lazilyAsync(create)Source.lazyFuture(create)延迟构造Future[T]lazySource的物化值为NotUsed迁移时把lazily替换为lazySource、把lazilyAsync替换为lazyFuture即可无需改动调用处与下游代码。lazySource还提供了更进阶的变体lazyFutureSource工厂返回Future[Source]见 Source.scala 第 612-613 行和对应的lazyCompletionStageSource以及 Flow/Sink 侧的Flow.lazyFlow、Sink.lazySink它们是同一“延迟物化”思路在流图不同位置的推广详见 lazySource.md 的 See also。八、小结Source.lazily的核心价值可归纳为三点延迟创建与物化工厂函数在首个需求到达时才执行内部 Source 的物化值通过Future/CompletionStage暴露给调用方每次物化一个全新实例工厂每次物化都会调用一次是安全持有可变状态的标准做法纯包装语义发射与完成行为完全取决于内部 Source。同时要记住它的两个边界异步边界/预取可能提前触发工厂下游在需求前取消时物化值会以NeverMaterializedException失败。由于 2.6.0 起已弃用生产代码请直接使用行为完全一致的Source.lazySource。赞分享后端并发编程异步编程【免费下载链接】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 Source.lazySource 详解延迟创建与物化的按需 SourceAkka Streams Source.lazySource 详解延迟创建与物化的按需 Source 导读 Source.lazySource 是 Akka后端并发编程异步编程Angular 延迟视图触发配置诊断 NG8021deferTriggerMisconfiguration深度解析Angular 延迟视图触发配置诊断 NG8021deferTriggerMisconfiguration深度解析 deferTriggerMisconfi后端并发编程异步编程Akka Streams Flow.lazyFlow 详解延迟创建与物化嵌套 FlowAkka Streams Flow.lazyFlow 详解延迟创建与物化嵌套 Flow Flow.lazyFlow 是 Akka Streams 中用于将 F后端并发编程异步编程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

相关推荐

三步跑通 Open-Meteo 天气 API 自托管:新手从零到上线的部署指南
三步跑通 Open-Meteo 天气 API 自托管:新手从零到上线的部署指南

三步跑通 Open-Meteo 天气 API 自托管:新手从零到上线的部署指南 【免费下载链接】open-meteo Free Weather Forecast API for non-commercial use 项目地址: https://gitcode.com/GitHub_Trending/op/open-meteo 如果你想在项目里免费查询全球天气预报、又不… · 2026/9/24 15:53:11

FerretDB 插入操作实战:insertOne 与 insertMany 的用法、响应解析与底层实现
FerretDB 插入操作实战:insertOne 与 insertMany 的用法、响应解析与底层实现

FerretDB 插入操作实战:insertOne 与 insertMany 的用法、响应解析与底层实现 【免费下载链接】FerretDB A truly Open Source MongoDB alternative 项目地址: https://gitcode.com/gh_mirrors/fe/FerretDB 本篇技术指南以 FerretDB 官方文档《Insert operat… · 2026/9/24 15:52:58

帝舵南京售后维修点丨门店位置、预约方式及咨询电话查询(2026年最新)
帝舵南京售后维修点丨门店位置、预约方式及咨询电话查询(2026年最新)

一份覆盖南京全域、同步全国标准的帝舵维保网点公示通告,帝舵南京售后维修点丨门店位置、预约方式及咨询电话查询(2026年最新),系统整合南京本地正规直营维保点位、精准出行方案、预约渠道与最新咨询热线,南京及江苏全域、皖东周边城市表主可… · 2026/9/24 15:52:58

Yii2 别名(Aliases)完全指南:从 `@` 符号到路径/URL 解析的底层机制
Yii2 别名(Aliases)完全指南:从 `@` 符号到路径/URL 解析的底层机制

后端Web框架 【免费下载链接】yii2 Yii 2: The Fast, Secure and Professional PHP Framework 项目地址: https://gitcode.com/gh_mirrors/yi/yii2 点击查看 免费下载 别名(Aliases)是 Yii 2 框架中表示文件路径和 URL 的轻量级符号机制&… · 2026/9/24 16:33:39

文化课教培数字化:课时自动核算 + 家校互动提升续费率完整方案
文化课教培数字化:课时自动核算 + 家校互动提升续费率完整方案

前言中小教培机构数字化转型,很多校长最先想到的功能是排课、消课,但在长期运营过程中,两个痛点会持续消耗机构大量人力成本:一是每月教师课时薪酬核算,二是老生续课留存。 尤其是文化课学科机构,课程类型复… · 2026/9/24 16:33:27

Sure 仓库的 AI 指令适配层:多 Harness 指令入口的统一维护实战
Sure 仓库的 AI 指令适配层:多 Harness 指令入口的统一维护实战

金融科技后端前端移动开发桌面应用AI 应用 【免费下载链接】sure The personal finance app for everyone (by everyone) 项目地址: https://gitcode.com/gh_mirrors/sure5/sure 点击查看 免费下载 本篇技术指南围绕 Sure(个人财务管理应用)… · 2026/9/24 16:33:20

CSP-S 2026 初赛试题解析(第二部分:阅读程序题(第一题))精讲
CSP-S 2026 初赛试题解析(第二部分:阅读程序题(第一题))精讲

2026 CSP-S 第一轮真题第二部分阅读程序第 1 题:《二进制除法》答案是:16:对✅️,17:对✅️,18:错❌️;19:C,20:B,21:C。一… · 2026/9/24 16:33:20

Visual C++ 6.0 MFC 单文档工程的多语言实现方案
Visual C++ 6.0 MFC 单文档工程的多语言实现方案

1. 引言 在 Visual C++ 6.0 时代,MFC(Microsoft Foundation Classes)是 Windows 桌面应用程序开发的主流框架。然而,其单文档工程(SDI)存在一个显著的限制:一个工程只能关联一个资源文件(.rc)。这意味着编译生成的可执行文件(.exe)默认只能包含一种语言的界面资源(… · 2026/9/24 16:33:14

四路can转4G在现场应用中有什么问题?
四路can转4G在现场应用中有什么问题?

一、现场使用 SG‑CAN‑4G‑410 网关,电脑通过网口配置设备,配置软件搜索不到设备,需要从哪些方面排查处理。 首先确认设备供电正常,PWR 电源灯常亮,RUN 系统指示灯处于闪烁运行状态;电脑网线连接设备 LAN … · 2026/9/24 16:33:14

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

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

1D-CNN时间序列建模实战:从Conv1d原理到工业落地
1D-CNN时间序列建模实战:从Conv1d原理到工业落地

简介:面向时间序列数据建模的一维卷积神经网络完整实现,适合深度学习入门者及需要快速验证时序模型的研究者,能够从音频、文本、传感器或股价等序列中挖掘局部特征与时间依赖。压缩包体积很小,只有3KB,内含3个Python脚… · 2026/9/24 0:00:26

柔软的L:汉语语流中被忽视的舌肌张力控制
柔软的L:汉语语流中被忽视的舌肌张力控制

1. 这个“L”不是字母表里的L,而是舌尖上的L最近在几个方言群和语音教学社群里,反复看到有人发一句:“也说字母L:柔软的长舌”。初看以为是英语发音课笔记,点开才发现全是方言爱好者、播音系学生、语言康复师甚至戏曲演… · 2026/9/24 0:00:44

了解更多?预约专属演示

我们的顾问将为您一对一讲解产品与方案

企业微信二维码