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

Akka Streams Source.lazyCompletionStage 详解:按需延迟创建单元素流的实现原理与实战用法

发布时间:2026/9/24 2:24:15 来源:云帆数科 栏目:资讯中心
Akka Streams Source.lazyCompletionStage 详解:按需延迟创建单元素流的实现原理与实战用法
后端并发编程异步编程【免费下载链接】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.lazyCompletionStage是 Akka Streams 中用于延迟创建的 Source 操作符它不会在流被构建或物化时立即执行用户提供的工厂函数而是等到下游真正产生需求demand时才调用工厂、取得一个CompletionStage并在该阶段成功完成后把唯一元素发射给下游随后流立即完成。本文以 lazyCompletionStage.md 为骨架结合 javadsl/Source.scala 与 scaladsl/Source.scala 的源码实现、以及LazyAndFutureSourcesTest与LazySourceSpec等测试用例带你掌握它的语义、Reactive Streams 契约、延迟失效的场景异步边界预取以及与Source.future、Source.lazySingle、Source.lazySource等近亲操作符的取舍。什么是 Source.lazyCompletionStage从文档的定义看Source.lazyCompletionStage的核心行为是Defers creation of a future of a single element source until there is demand. 把单个元素 source 对应的 future的创建延迟到存在需求时。具体来说它做三件事延迟调用工厂用户提供的工厂函数factory在流构建与物化阶段都不会被执行只有当下游第一个需求demand到达时才被调用发射唯一元素工厂返回的CompletionStage成功完成后其值作为唯一一个流元素发射给下游失败传播如果工厂函数本身抛异常或返回的CompletionStage以失败结束整个流随之失败。文档同时给出了 Reactive Streams 语义emits发射当存在下游需求且元素工厂返回的 future 已完成时发射唯一元素completes完成发射完这唯一一个元素后流立即完成。与相关操作符的关系lazyCompletionStage属于 Akka Streams 的延迟创建操作符家族在 javadsl/Source.scala 与 scaladsl/Source.scala 中与下列操作符相邻定义便于对比操作符工厂延迟调用工厂返回值下游看到的内容lazySingle是普通值T单个元素lazyCompletionStage即 Scala 侧lazyFuture是CompletionStage[T]/Future[T]单个元素lazySource是一个Source该 Source 的全部元素lazyCompletionStageSource即 Scala 侧lazyFutureSource是CompletionStage[Source]/Future[Source]内层 Source 的全部元素future/completionStage否创建时即已给定已存在的Future/CompletionStage单个元素可以看出lazyCompletionStage处于按需创建与异步单元素两条语义的交汇点它既像lazySingle一样延迟执行工厂又像completionStage一样把异步结果作为唯一元素发射。源码实现它是如何做到延迟的在 Scala DSL 中Source.lazyFutureJava DSL 的lazyCompletionStage直接委托给它的实现异常简洁def lazyFutureT Future[T]): Source[T, NotUsed] single(()).mapAsyncUnordered(1)(_ create()).withAttributes(DefaultAttributes.lazyFuture)见 scaladsl/Source.scala。这段实现可以拆解为三个关键点single(())先用一个立即发射Unit的单元素 Source 打底。由于它只有一个元素且要等下游 demand 才会发射因此下游的第一次需求会先抵达这里mapAsyncUnordered(1)并发度为 1 的异步映射。只有当上游即single(())把Unit元素发下来时create()工厂才被真正执行返回的Future完成时其值被发射出去。这既保证了工厂只被调用一次又保证了发射时机是 future 完成之后withAttributes(DefaultAttributes.lazyFuture)为该操作符挂上专属属性name 等便于调试与监控。而 Java DSL 侧 javadsl/Source.scala 只是做了一层适配public T SourceT, NotUsed lazyCompletionStage(CreatorCompletionStageT create) { return scaladsl.Source.lazyFuture(() - create.create().asScala).asJava; }即把 Java 的CreatorCompletionStageT转成 Scala 的() Future[T]复用同一套 Scala 实现因此两者的语义完全一致。为什么用 mapAsyncUnordered 而不是普通 map使用mapAsyncUnordered(1)而非map或flatMapConcat的深层原因在于工厂返回的是一个异步Future/CompletionStage需要以异步方式等待其结果并发度设为 1保证同一时刻最多只有一个工厂调用在进行杜绝了重复调用一旦Future失败mapAsyncUnordered会把该失败作为流失败向下游传播——这正是文档中future 失败则流失败语义的实现机制。何时真正调用工厂从 GraphStage 到 FutureSource对于延迟创建的语义lazyFuture借助singlemapAsyncUnordered的链式结构即可实现。而家族中更复杂的lazySource/lazyFutureSource则使用了专门的GraphStageimpl/LazySource.scala 中的LazySource[T, M]处理延迟到有需求时才物化内层 Sourceimpl/fusing/GraphStages.scala 中的FutureFlattenSource[T, M]处理等待 future 完成后展开为内层 Source其内部通过FutureSource同文件 GraphStages.scala等待异步结果。对比可见lazyCompletionStage是这条延迟创建链路里最轻量的一环——它只需等一个异步值而不需要等一个完整的内层流因此实现上没有引入新的 GraphStage而是直接复用singlemapAsyncUnordered的组合这本身也印证了其发射单元素后立即完成的语义。实战用法Java 与 Scala 双侧示例Java 示例来自 LazyAndFutureSourcesTest.java 的测试用例展示了最典型的用法——工厂返回一个已完成或任意时刻完成的CompletableFutureimport akka.NotUsed; import akka.actor.ActorSystem; import akka.stream.javadsl.Source; import akka.stream.javadsl.Sink; import java.util.Arrays; import java.util.List; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletionStage; ActorSystem system ActorSystem.create(lazy-completion-stage-example); // 工厂不会在下面这行执行只会在有下游 demand 时才执行 SourceString, NotUsed src Source.lazyCompletionStage(() - CompletableFuture.completedFuture(one)); CompletionStageListString result src.runWith(Sink.seq(), system); // 断言结果恰好是 [one] assertEquals(Arrays.asList(one), result.toCompletableFuture().get(3, TimeUnit.SECONDS));关键点Source.lazyCompletionStage(...)这一行本身不会触发工厂只有runWith(Sink.seq(), system)物化并产生需求后工厂才被调用。Scala 示例Scala DSL 侧对应的方法是Source.lazyFuture。参考 LazySourceSpec.scala 中的测试可以覆盖以下场景import akka.actor.ActorSystem import akka.stream.scaladsl.{Sink, Source} import scala.concurrent.{Future, Promise} implicit val system ActorSystem(lazy-future-example) implicit val ec system.dispatcher // 场景 1工厂返回已完成的 Future正常发射单元素 val seq1 Source.lazyFuture(() Future.successful(1)).runWith(Sink.seq) // seq1.futureValue Seq(1) // 场景 2工厂返回尚未完成的 Future等它完成后再发射 val promise Promise[Int]() val seq2 Source.lazyFuture(() promise.future).runWith(Sink.seq) promise.success(1) // seq2.futureValue Seq(1) // 场景 3下游没有任何需求这里直接取消工厂绝不被调用 val constructed new java.util.concurrent.atomic.AtomicBoolean(false) Source .lazyFuture { () constructed.set(true) Future.successful(1) } .watchTermination()(Keep.right) .toMat(Sink.cancelled)(Keep.left) .run() // constructed.get() false —— 证明没有 demand 就没有调用第三个场景直接验证了本文开头文档中延迟创建的核心承诺即使流已经被物化run()已执行只要下游没有需求这里使用Sink.cancelled工厂就不会被调用。失败语义工厂异常与 Future 失败文档明确指出If the future or the factory fails the stream is failed.如果 future 或工厂失败流就失败。LazySourceSpec.scala 用三个测试覆盖了全部失败路径失败场景测试代码结果工厂函数直接抛异常Source.lazyFuture(() throw failure)流以failure失败工厂返回一个已失败的 FutureSource.lazyFuture(() Future.failed(failure))流以failure失败工厂返回的 Future 后来才失败val p Promise[Int](); Source.lazyFuture(() p.future); p.failure(failure)流以failure失败三个场景统一断言termination.failed.futureValue failure即通过watchTermination()观察到的终止信号都是同一个失败原因。这印证了无论是调用时刻还是异步完成时刻发生失败最终都以流失败的形式呈现给下游不存在静默吞掉异常的情况。重要注意事项异步边界会破坏惰性文档特别强调了一个容易踩坑的行为Note that asynchronous boundaries (and other operators) in the stream may do pre-fetching which counter acts the laziness and will trigger the factory immediately. 流中的异步边界以及其他操作符可能会进行预取pre-fetching这会抵消惰性并立即触发工厂。这是因为 Akka Streams 基于 Reactive Streams 实现背压异步边界两侧通过独立的执行上下文和内部缓冲区通信默认会预先请求一批元素prefetch以提升吞吐。一旦预取发生下游的需求就已产生lazyCompletionStage的工厂自然会被提前调用惰性失效。实用建议如果下游紧跟着async边界、buffer等带缓冲/预取语义的操作符不要假设工厂一定会被推迟到最后时刻若需要真正严格的直到最后一刻才执行应把lazyCompletionStage放在流水线最末端、紧贴真正消费元素的 Sink 之前或通过控制预取参数来调整行为该操作符适合代价高昂、希望按需创建的场景例如连接外部服务前的认证 token 获取、按需加载配置、需要等待异步准备完成后再产出数据的边界场景。与 Source.completionStage 的对比何时用哪个同样发射单个异步元素Source.completionStageScala 侧Source.future与Source.lazyCompletionStage的差别在于异步结果的来源时机Source.completionStage(stage)CompletionStage在构建流时就已经存在通常已在某处被启动操作符只是把它接入流中即使没有任何下游需求这个异步任务也已经开始执行Source.lazyCompletionStage(() - stage)CompletionStage是在收到下游需求那一刻才由工厂创建并启动的没有需求就完全没有副作用产生。参考 javadsl/Source.scala 中completionStage的实现它直接调用future(completionStage.asScala)即传入的CompletionStage是现成的工厂调用与否无从谈起。因此选择依据很明确需要惰性副作用按需发生选lazyCompletionStage异步结果已经存在、只需要接入流中选completionStage。家族全览从 lazySingle 到 lazyCompletionStageSourcelazyCompletionStage只是 Akka Streams 延迟创建家族的一员。为便于在实际工程中做出正确选择这里把 scaladsl/Source.scala 中全部相关方法整理成速查表方法工厂签名物化值下游语义lazySingle(create: () T)同步返回单值NotUsed有需求时调用发射单元素后完成lazyFuture(create: () Future[T])异步返回单值NotUsed有需求时调用Future 完成后发射单元素然后完成即文档主题lazySource(create: () Source[T, M])返回内层 SourceFuture[M]内层物化值有需求时物化内层 Source其行为完全等同于直接使用该 SourcelazyFutureSource(create: () Future[Source[T, M]])异步返回内层 SourceFuture[M]有需求时调用等 Future 完成后展开内层 Source其中lazySource与lazyFutureSource还具备更精细的物化值语义如果下游在工厂被调用前就取消或失败物化值会被以akka.stream.NeverMaterializedException失败见 scaladsl/Source.scala 与 javadsl/Source.scala 的文档注释。这一设计保证了从未被物化的内层流不会产生悬空的物化值相关行为也在 LazySourceSpec.scala 中有测试覆盖断言lazyFutureSourceMatval.failed.futureValue shouldBe a[NeverMaterializedException]。小结Source.lazyCompletionStage通过single(())mapAsyncUnordered(1)的轻量组合实现了下游有需求才调用工厂、Future 完成后发射唯一元素、失败即流失败的完整语义是 Akka Streams 中按需异步生产的首选操作符。使用它时务必记住两点没有下游需求就没有任何副作用可用Sink.cancelled验证以及异步边界与预取会破坏惰性。在需要严格按需启动异步任务、避免无谓资源消耗的场景下它比Source.completionStage更具表达力而当异步结果早已存在时则应优先选择更简单的completionStage。延伸阅读操作符索引Source operatorsJava DSL 实现javadsl/Source.scalaScala DSL 实现scaladsl/Source.scala底层 GraphStageimpl/LazySource.scala、impl/fusing/GraphStages.scalaJava 测试LazyAndFutureSourcesTest.javaScala 测试LazySourceSpec.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 的 Source.lazySingle 操作符按需延迟创建单元素流的完整指南Akka Streams 的 Source.lazySingle 操作符按需延迟创建单元素流的完整指南 Source.lazySingle 是 Akka St后端并发编程异步编程Akka Streams Source.lazySource 详解延迟创建与物化的按需 SourceAkka Streams Source.lazySource 详解延迟创建与物化的按需 Source 导读 Source.lazySource 是 Akka后端并发编程异步编程深入 Akka Streams 的 Source.lazyCompletionStageSource按需延迟创建「未来 Source」深入 Akka Streams 的 Source.lazyCompletionStageSource 按需延迟创建「未来 Source」 Source.laz后端并发编程异步编程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

相关推荐

ESP32+3D打印瓦力机器人DIY教程:从选型到调试
ESP32+3D打印瓦力机器人DIY教程:从选型到调试

/* 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 2:24:15

Unity UGUI DropDown问题
Unity UGUI DropDown问题

复制出来的下拉菜单比模板短预制体里面Content加Vertical Layout Group,Content Size Fitter,然后模板激活一下让布局生效,再隐藏保存。两个自动布局组件这样设: · 2026/9/24 2:24:08

XY2-100协议详解:从差分信号到FPGA/MCU振镜控制方案
XY2-100协议详解:从差分信号到FPGA/MCU振镜控制方案

/* 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 2:24:08

Nginx UI 开发环境搭建:基于 Devcontainer 的一键容器化开发与多节点集群调试指南
Nginx UI 开发环境搭建:基于 Devcontainer 的一键容器化开发与多节点集群调试指南

后端前端运维MCP 服务 【免费下载链接】nginx-ui Yet another WebUI for Nginx 项目地址: https://gitcode.com/gh_mirrors/ngi/nginx-ui 点击查看 免费下载 导读 本文基于 Nginx UI 仓库的 docs/guide/devcontainer.md 与 .devcontainer 目录下的真实配置&#x… · 2026/9/24 3:02:36

Orleans 生产环境部署与运维完全指南:集群规划、平台选型与故障恢复
Orleans 生产环境部署与运维完全指南:集群规划、平台选型与故障恢复

后端微服务 【免费下载链接】orleans Cloud Native application framework for .NET 项目地址: https://gitcode.com/gh_mirrors/or/orleans 点击查看 免费下载 导读 本文是 Orleans 生产部署与运维的完整操作指南。Orleans 的生产形态是一组通过 TCP 直连的 silo… · 2026/9/24 3:02:11

深入解析 wandb core 中的 Go JOSE v4:基于 RFC 7515/7516/7519 的 JWS、JWE 与 JWT 实现指南
深入解析 wandb core 中的 Go JOSE v4:基于 RFC 7515/7516/7519 的 JWS、JWE 与 JWT 实现指南

机器学习深度学习数据可视化可观测性 【免费下载链接】wandb The AI developer platform. Use Weights & Biases to train and fine-tune models, and manage models from experimentation to production. 项目地址: https://gitcode.com/gh_mirrors/wa/wandb 点… · 2026/9/24 3:02:05

多轨道二次编辑怎么用
多轨道二次编辑怎么用

多轨道二次编辑是剪映专业版针对初步剪辑完成的AI生成内容做精修的方法:你可以在已经排好的时间线上,只针对不满意的单个AI片段单独发起二次生成替换,保留其他轨道的内容和整体剪辑结构不变,不用重新调整整个成片的编排。这种方式… · 2026/9/24 3:01:59

Kornia 迁移指南:BoxMotTracker 移除与基于 boxmot + RTDETRDetectorBuilder 的替代方案
Kornia 迁移指南:BoxMotTracker 移除与基于 boxmot + RTDETRDetectorBuilder 的替代方案

计算机视觉深度学习人工智能图像处理 【免费下载链接】kornia 🐍 空间人工智能的几何计算机视觉库 项目地址: https://gitcode.com/kornia/kornia 点击查看 免费下载 本篇技术指南聚焦 Kornia 开源仓库中的一项破坏性变更(Migration 004&… · 2026/9/24 3:01:59

Sliver 网络侦察命令组实战:ifconfig 与 netstat 的架构、实现与使用详解
Sliver 网络侦察命令组实战:ifconfig 与 netstat 的架构、实现与使用详解

网络安全 【免费下载链接】sliver Adversary Emulation Framework 项目地址: https://gitcode.com/gh_mirrors/sl/sliver 点击查看 免费下载 导读 本篇技术指南以 Sliver 客户端 client/command/network 命令组为主线,深入解析其两个核心网络侦察命令 … · 2026/9/24 3:01:47

基于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

了解更多?预约专属演示

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

企业微信二维码