news 2026/9/23 11:25:15

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

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Akka Streams Source.lazyFuture 操作符完全指南:延迟创建单元素 Future 的惰性数据源
  • 后端
  • 并发编程
  • 异步编程

【免费下载链接】akka-core

A 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.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):

  1. single(()):先构造一个发射单个Unit元素的 Source,作为"触发器";
  2. mapAsyncUnordered(1):以并行度 1 的方式对触发器元素调用create(),得到Future[T],并在其完成时把结果T发射给下游;
  3. 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]延迟物化一个完整 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(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覆盖了五类场景,直接印证了本文档的全部语义:

  1. happy path(Future 已成功)Source.lazyFuture(() => Future.successful(1)).runWith(Sink.seq)结果为Seq(1)
  2. happy path(Future 稍后完成):工厂返回Promise的 Future,promise.success(1)后结果同样为Seq(1)
  3. 无需求不构造:用AtomicBoolean标记工厂是否执行,配合Sink.cancelled取消下游后,constructed.get()false,且流正常终止;
  4. 工厂函数抛异常() => throw failure使流以该异常失败;
  5. 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.

项目地址:https://gitcode.com/gh_mirrors/ak/akka-core
点击查看免费下载

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/23 11:23:20

硬件测试规范实战:从原理图审查到自动化脚本的完整指南

简介&#xff1a;这份硬件测试方案文档面向硬件工程师、测试人员及电子相关专业学生&#xff0c;聚焦整机与单板两类测试场景&#xff0c;帮助读者建立从测试目的、适用范围到判定准则的完整测试框架。资源包内含1个doc文件&#xff0c;约7.95MB&#xff0c;共76页&#xff0c;…

作者头像 李华
网站建设 2026/9/23 11:21:25

PC与Android视频传输方案全解析

1. 多平台视频传输需求背景在跨设备协作成为常态的今天&#xff0c;PC与移动端之间的视频传输已成为刚需。无论是剪辑师需要将成品快速发送到手机预览&#xff0c;还是普通用户想在大屏设备上观看下载影片&#xff0c;高效的文件传输方案都能显著提升工作效率和使用体验。Andro…

作者头像 李华
网站建设 2026/9/23 11:20:25

公益src一次简单的验证码绕过、弱口令

1、在漏洞平台&#xff0c;公益SRC上&#xff0c;找一个网站&#xff0c;找到登录处2、抓包&#xff0c;发现密码明文&#xff0c;放到Repeater&#xff0c;多次G&#xff4f;&#xff0c;发现没有验证码&#xff0c;也不限制次数&#xff13;、进行爆破&#xff0c;得到密码&a…

作者头像 李华