- 后端
- 并发编程
- 异步编程
【免费下载链接】akka-core
A platform to build and run apps that are elastic, agile, and resilient. SDK, libraries, and hosted environments.
本文围绕 Akka 项目中Source.lazyFuture操作符展开:它把"创建一个单元素 Future"这件事推迟到下游真正出现需求(demand)时才执行,从而避免无关的副作用与资源浪费,是构建按需计算、按需加载数据源的实用工具。读完本文,你将掌握lazyFuture的签名与语义、底层实现原理、与lazySingle/lazySource/lazyFutureSource家族操作符的差异,以及它在 Scala 与 Java 两种 DSL 下的实战用法与边界条件。
概览:什么是 Source.lazyFuture
Source.lazyFuture是 Akka Streams 提供的一个"惰性"(lazy)Source 工厂方法。与普通 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.scala):
single(()):先构造一个发射单个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) | 普通值 | 延迟计算一个同步元素 | NotUsed |
lazyFutureT => Future[T]) | Future[T] | 延迟创建一个异步元素(本文主角) | NotUsed |
lazySourceT, M => Source[T, M]) | Source[T, M] | 延迟物化一个完整 Source | Future[M] |
lazyFutureSourceT, M => Future[Source[T, M]]) | Future[Source[T, M]] | 延迟创建一个 Future 包裹的 Source | Future[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(s"constructed = ${constructed.get()}") // 输出 false }失败传播
三种失败途径都会让整个流失败,且携带原始异常:
// 1) 工厂函数直接抛异常 Source.lazyFuture(() => throw new RuntimeException("couldn't 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-L353):
import 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"); Source<Integer, NotUsed> src = Source.lazyCompletionStage( (Creator<CompletionStage<Integer>>) () -> 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 path(Future 已成功):
Source.lazyFuture(() => Future.successful(1)).runWith(Sink.seq)结果为Seq(1); - happy path(Future 稍后完成):工厂返回
Promise的 Future,promise.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-core
A platform to build and run apps that are elastic, agile, and resilient. SDK, libraries, and hosted environments.
相关推荐
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),仅供参考