十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Akka Streams Source.lazyFuture 操作符完全指南:延迟创建单元素 Future 的惰性数据源

Akka Streams Source.lazyFuture 操作符完全指南:延迟创建单元素 Future 的惰性数据源 后端并发编程异步编程【免费下载链接】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 项目中Source.lazyFuture操作符展开它把创建一个单元素 Future这件事推迟到下游真正出现需求demand时才执行从而避免无关的副作用与资源浪费是构建按需计算、按需加载数据源的实用工具。读完本文你将掌握lazyFuture的签名与语义、底层实现原理、与lazySingle/lazySource/lazyFutureSource家族操作符的差异以及它在 Scala 与 Java 两种 DSL 下的实战用法与边界条件。概览什么是 Source.lazyFutureSource.lazyFuture是 Akka Streams 提供的一个惰性lazySource 工厂方法。与普通 Source 在物化materialization时立即创建元素不同lazyFuture将用户提供的工厂函数create的调用推迟到下游第一个需求demand到达之时当返回的 Future成功完成时其结果作为单个流元素向下游发射如果 Future 失败或工厂函数本身抛出异常整个流以该异常失败fail发射完这唯一一个元素后流立即正常完成complete。该操作符在 akka-docs 官方文档 中归属于 Source 操作符refSource operators对应的 Reactive Streams 语义为语义说明emits当下游存在需求且元素工厂返回的 Future 已完成时completes在发射完这唯一一个元素之后签名与类型lazyFuture在 Scala DSL 中的完整签名为def lazyFutureT Future[T]): Source[T, NotUsed]create返回Future[T]的工厂函数签名是() Future[T]返回值Source[T, NotUsed]即发射类型为T、物化值为NotUsed不产生有意义的物化值的 Source。该签名定义在 akka-stream/src/main/scala/akka/stream/scaladsl/Source.scala#L573-L574其官方 API 文档入口为 apidocSource.lazyFuture。底层实现原理lazyFuture的实现非常精巧——它不是独立的 GraphStage而是由两个既有操作组合而成def lazyFutureT Future[T]): Source[T, NotUsed] single(()).mapAsyncUnordered(1)(_ create()).withAttributes(DefaultAttributes.lazyFuture)实现要点对应 Source.scalasingle(())先构造一个发射单个Unit元素的 Source作为触发器mapAsyncUnordered(1)以并行度 1 的方式对触发器元素调用create()得到Future[T]并在其完成时把结果T发射给下游withAttributes(DefaultAttributes.lazyFuture)为操作符打上lazyFuture的默认属性名便于日志与调试见 akka-stream/src/main/scala/akka/stream/impl/Stages.scala#L149-L151。正是因为外层是single(())只有当下游真正产生需求、该单元素被请求时mapAsyncUnordered才会执行create()。若下游从不拉取例如使用Sink.cancelled立即取消工厂函数永远不会被调用。与 lazy 家族操作符的关系lazyFuture不是孤立存在的。在 Source.scala 中它属于一个完整的延迟创建操作符家族四者按创建对象的粒度递进操作符工厂返回类型延迟创建的内容物化值lazySingleT T)普通值延迟计算一个同步元素NotUsedlazyFutureT Future[T])Future[T]延迟创建一个异步元素本文主角NotUsedlazySourceT, M Source[T, M])Source[T, M]延迟物化一个完整 SourceFuture[M]lazyFutureSourceT, M Future[Source[T, M]])Future[Source[T, M]]延迟创建一个 Future 包裹的 SourceFuture[M]其中lazySingle是同步版single(()).map(_ create())lazyFuture是异步版single(()).mapAsyncUnordered(1)(_ create())lazySource/lazyFutureSource则基于独立的LazySourceGraphStage 实现akka-stream/src/main/scala/akka/stream/impl/LazySource.scala可发射多个元素且其物化值通过Promise在内部 Source 物化时完成若下游在工厂被调用前取消物化值会以NeverMaterializedException失败。选型建议只需要发射单个结果、且结果来自异步计算时用lazyFuture结果可以同步算出时用lazySingle需要发射多个元素或要拿到内部 Source 的物化值时升级到lazySource/lazyFutureSource。注意惰性并非绝对官方文档特别强调了一个关键限制见 lazyFuture.md流中的异步边界asynchronous boundaries和其他操作符可能做预取pre-fetching这会抵消惰性导致工厂函数被立即触发。也就是说如果lazyFuture后面接了会提前向下游拉取的操作如buffer、异步边界、带缓冲的算子下游需求可能在物化后很快到达甚至在下游真正想要数据之前就已触发create()。在需要严格保证绝不在需求出现前执行副作用的场景中应避免在lazyFuture与消费者之间放置预取型算子。实战示例Scala DSL以下示例可在 Akka Streams 2.x 的 Scala 工程中直接运行。基本用法Future 已就绪import akka.actor.ActorSystem import akka.stream.scaladsl.{ Sink, Source } implicit val system: ActorSystem ActorSystem(lazyFuture-demo) import system.dispatcher val seq Source.lazyFuture(() Future.successful(1)).runWith(Sink.seq) // seq 完成后结果为 Seq(1)发射单个元素后流即完成延迟到 Promise 完成工厂函数返回的 Future 可以稍后才完成流会一直等待其完成后再发射import scala.concurrent.Promise val promise Promise[Int]() val seq Source.lazyFuture(() promise.future).runWith(Sink.seq) promise.success(1) // 稍后完成 seq.foreach(println) // 输出 Seq(1)无需求时不构造这是lazyFuture的核心价值下游不拉取工厂就绝不执行import java.util.concurrent.atomic.AtomicBoolean val constructed new AtomicBoolean(false) val termination Source .lazyFuture { () constructed.set(true) Future.successful(1) } .watchTermination()(Keep.right) .toMat(Sink.cancelled)(Keep.left) // 下游立即取消不产生需求 .run() termination.foreach { _ println(sconstructed ${constructed.get()}) // 输出 false }失败传播三种失败途径都会让整个流失败且携带原始异常// 1) 工厂函数直接抛异常 Source.lazyFuture(() throw new RuntimeException(couldnt create)) // 2) 工厂返回已失败的 Future Source.lazyFuture(() Future.failed(new RuntimeException(future failed))) // 3) 工厂返回的 Future 之后失败 val p Promise[Int]() Source.lazyFuture(() p.future) p.failure(new RuntimeException(later failure))以上三种情形下游都会收到对应的失败信号流终止。实战示例Java DSL在 Java DSL 中lazyFuture的对应方法是Source.lazyCompletionStage它内部把CompletionStage适配为 ScalaFuture后委托给lazyFuture见 akka-stream/src/main/scala/akka/stream/javadsl/Source.scala#L350-L353import akka.actor.ActorSystem; import akka.japi.function.Creator; import akka.stream.javadsl.Sink; import akka.stream.javadsl.Source; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletionStage; ActorSystem system ActorSystem.create(lazyFuture-demo); SourceInteger, NotUsed src Source.lazyCompletionStage( (CreatorCompletionStageInteger) () - CompletableFuture.completedFuture(42)); src.runWith(Sink.seq(), system) .thenAccept(seq - System.out.println(seq)); // [42]Java DSL 中同家族还包括lazySingle同步值与lazySource返回CompletionStage[M]物化值详见 javadsl/Source.scala。测试用例对语义的验证仓库中的 LazySourceSpec.scala 对Source.lazyFuture覆盖了五类场景直接印证了本文档的全部语义happy pathFuture 已成功Source.lazyFuture(() Future.successful(1)).runWith(Sink.seq)结果为Seq(1)happy pathFuture 稍后完成工厂返回Promise的 Futurepromise.success(1)后结果同样为Seq(1)无需求不构造用AtomicBoolean标记工厂是否执行配合Sink.cancelled取消下游后constructed.get()为false且流正常终止工厂函数抛异常() throw failure使流以该异常失败Future 失败Future.failed(failure)或Promise稍后failure(failure)流均以该异常失败。这些用例可从测试入口 akka-stream-tests/src/test/scala/akka/stream/scaladsl/LazySourceSpec.scala 查看完整实现。典型应用场景与小结Source.lazyFuture适合以下场景按需执行开销较大的初始化例如仅在消费者真正需要时才发起远程调用、读取数据库或执行计算避免应用启动阶段触发无关副作用延迟错误把可能抛异常的代码包进工厂函数将错误从物化阶段推迟到需求阶段交由流的失败信号统一处理串联异步单值与mapAsync等算子配合构造先等待、后单值发射的数据源。同时务必牢记两点边界其一流的预取与异步边界可能提前触发工厂无法保证绝对的惰性其二lazyFuture只发射一个元素需要多元素或完整 Source 语义时应转向lazySource/lazyFutureSource。理解这些行为后你就能在 Akka Streams 中精准地驾驭延迟数据源的创建时机。赞分享后端并发编程异步编程【免费下载链接】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 Flow.futureFlow 操作符延迟创建内部流与按需物化的完整指南Akka Streams Flow.futureFlow 操作符延迟创建内部流与按需物化的完整指南 导读 Flow.futureFlow 是 Akka Str后端并发编程异步编程Akka Streams groupedWeighted 操作符完全指南按元素权重聚合流Akka Streams groupedWeighted 操作符完全指南按元素权重聚合流 groupedWeighted 是 Akka Streams 中用于后端并发编程异步编程Akka Streams delayWith 操作符详解按元素动态控制延迟的定时驱动流处理Akka Streams delayWith 操作符详解按元素动态控制延迟的定时驱动流处理 导读 delayWith 是 Akka Streams 中一类特殊后端并发编程异步编程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表