Akka Streams 的 Source.unfoldAsync 详解:基于 Future/CompletionStage 的状态驱动异步数据源
后端并发编程异步编程【免费下载链接】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.unfoldAsync是 Akka Streams 中Source家族的关键操作符之一它与同步的unfold行为一致但折叠函数返回的是 scala[Future] java[CompletionStage]因此非常适合用来实现以异步方式从外部服务Actor、数据库、HTTP 接口、文件系统按偏移量分批拉取数据这类有状态、有终止条件的数据源。读完本文你将掌握unfoldAsync的签名与状态机模型、Reactive Streams 语义、内部 GraphStage 实现原理以及用 Actor ask 模式实现分块读取的完整可运行示例。概述与unfold的异同unfoldAsync与同步的unfold行为完全一致——每次调用折叠函数时函数接收上一个状态并返回 scala[Option[(S, E)]] java[OptionalPairS, E]其中元组的第一个元素是传给下一次调用的新状态第二个元素是向下游发射的元素。唯一的区别在于Just likeunfoldbut the fold function returns a scala[Future] java[CompletionStage] which will cause the source to complete or emit when it completes.也就是说折叠函数返回的是异步计算结果Source 会等待该异步结果完成后再决定是发射元素还是完成流。这使得unfoldAsync可以直接对接一切以回调/异步返回结果的数据源例如向 Actor 发送 ask 请求并等待回复向数据库、对象存储或 HTTP 服务发起异步查询基于Future/CompletionStage包装的阻塞式 I/O。官方文档还明确指出它可以用来实现许多有状态的 Source而无需触碰更低层的GraphStageAPI即 stream-customize.md 中介绍的自定义流阶段编写方式。如果你的数据源是阻塞式的如同步网络或文件系统 API则文档建议优先考虑同步变体unfold的配套操作符详见 Source.unfold 文档 中对unfoldResource的说明。签名与 API 对比Scala API定义于 scaladsl/Source.scaladef unfoldAsyncS, E(f: S Future[Option[(S, E)]]): Source[E, NotUsed]s: S初始状态zero 状态折叠函数的第一次调用会收到它f: S Future[Option[(S, E)]]异步折叠函数返回一个在将来解析为Option的Future返回值Source[E, NotUsed]发射类型E的元素材质化值为NotUsed无有用材质化值。Java API定义于 javadsl/Source.scaladef unfoldAsyncS, E: Source[E, NotUsed]Java 版本使用CompletionStageOptionalPairS, E其中Pair来自akka.japi.Pair。与同步unfold的对照维度unfoldunfoldAsync折叠函数返回值scala[Option[(S, E)]] java[OptionalPairS,E]scala[Future[Option[(S, E)]]] java[CompletionStageOptionalPairS,E]何时发射/完成函数同步返回后立即决定异步结果完成后决定适用场景纯内存计算、同步迭代异步 I/O、Actor ask、外部服务调用状态推进函数返回的S作为下一次调用的状态同上但发生在 Future 完成之后从源码看Scala 的unfoldAsync直接基于内部的UnfoldAsyncGraphStage 构建Source.fromGraph(new UnfoldAsync(s, f))而 Java 版本则使用专门为CompletionStage优化的UnfoldAsyncJava实现impl/Unfold.scala两者共享同一个默认属性名unfoldAsync见 impl/Stages.scala。状态机模型与执行流程unfoldAsync是一个典型的有状态、按需demand-driven的数据源。它的运行可以抽象为如下循环以初始状态s0调用折叠函数f(s0)得到Future[Option[(S, E)]]当下游产生需求pull时等待该 Future 完成若结果为Some((newState, elem))向下游发射elem并把状态更新为newState然后等待下一次 pull 再调用f(newState)若结果为None完成complete整个流若 Future 失败以该异常使流失败fail。重复步骤 2直到返回None或流被下游取消。内部实现原理UnfoldAsyncGraphStageunfoldAsync的底层实现位于 impl/Unfold.scala核心代码展示了它是如何安全地把异步结果投递回流的执行线程的override def preStart(): Unit { asyncHandler getAsyncCallback[Try[Option[(S, E)]]](handle).invoke } private def handle(result: Try[Option[(S, E)]]): Unit result match { case Success(Some((newS, elem))) push(out, elem) state newS case Success(None) complete(out) case Failure(ex) fail(out, ex) } def onPull(): Unit { val future f(state) future.value match { case Some(value) handle(value) // 已完成的 Future立即处理 case None future.onComplete(asyncHandler)(ExecutionContext.parasitic) // 未完成注册回调 } }几个值得注意的实现细节getAsyncCallback机制Future 完成回调可能运行在任意线程如 Actor 派发线程、EC 线程池而 Akka Streams 规定只有流的执行线程才能安全地push元素。getAsyncCallback会把异步回调搬运回流执行线程这正是unfoldAsync能安全对接Future的关键已完成 Future 的快速路径future.value match { case Some(value) handle(value) }说明如果传入的 Future 已经完成则直接同步处理避免不必要的回调开销失败传播Future 失败会调用fail(out, ex)使整个流以该异常失败——因此务必确保折叠函数返回的 Future 能够妥善处理内部异常例如在map中处理边界条件而不是让 Future 意外失败状态更新时机状态只在Success(Some(...))时更新None或失败都不会推进状态这保证了状态流转的确定性。Java 专用实现UnfoldAsyncJavaimpl/Unfold.scala逻辑一致只是针对CompletableFuture的isDone/getNow/handle做了等价优化并约定Optional.empty()表示流结束、Pair.first为新状态、Pair.second为发射元素。Reactive Streams 语义unfoldAsync遵循标准背压back-pressure语义官方文档明确给出了两条核心语义emits发射当存在下游需求且 unfold 状态返回的 Future 完成并携带某个值时发射该值completes完成当 unfold 函数返回的 Future 完成且值为空scala[None] java[Optional.empty]时完成流。结合源码可以进一步明确发射和完成都发生在 Future 完成之后且都受下游需求驱动——没有下游 pull 就不会调用折叠函数onPull才触发f(state)因此unfoldAsync天然具备惰性拉取lazy pull特性不会在无人消费时提前执行异步请求。实战示例通过 Actor ask 分块读取数据官方文档unfoldAsync.md给出的示例场景是向一个模拟 Actor 按偏移量请求字节块Actor 返回Chunk消息当请求的偏移量超过数据末尾时返回空ByteString。我们希望通过unfoldAsync把它表示成一个在到达末尾时自然完成的ByteString流。这里的技巧是把 offset 作为每次调用之间传递的状态。定义 Actor 协议Scala完整代码见 UnfoldAsync.scalaobject DataActor { sealed trait Command case class FetchChunk(offset: Long, replyTo: ActorRef[Chunk]) extends Command case class Chunk(bytes: ByteString) }Java完整代码见 UnfoldAsync.javaclass DataActor { interface Command {} static final class FetchChunk implements Command { public final long offset; public final ActorRefChunk replyTo; public FetchChunk(long offset, ActorRefChunk replyTo) { this.offset offset; this.replyTo replyTo; } } static final class Chunk { public final ByteString bytes; public Chunk(ByteString bytes) { this.bytes bytes; } } }协议要点客户端携带offset发起请求Actor 回复Chunk如果请求的 offset 超出了数据末尾Actor 返回空的ByteString这是流的终止信号。用 unfoldAsync 实现分块拉取Scala 示例来自 UnfoldAsync.scala// actor we can query for data with an offset val dataActor: ActorRef[DataActor.Command] ??? import system.executionContext implicit val askTimeout: Timeout 3.seconds val startOffset 0L val byteSource: Source[ByteString, NotUsed] Source.unfoldAsync(startOffset) { currentOffset // ask for next chunk val nextChunkFuture: Future[DataActor.Chunk] dataActor.ask(DataActor.FetchChunk(currentOffset, _)) nextChunkFuture.map { chunk val bytes chunk.bytes if (bytes.isEmpty) None // end of data else Some((currentOffset bytes.length, bytes)) } }Java 示例来自 UnfoldAsync.javaActorRefDataActor.Command dataActor null; // lets say we got it from somewhere Duration askTimeout Duration.ofSeconds(3); long startOffset 0L; SourceByteString, NotUsed byteSource Source.unfoldAsync( startOffset, currentOffset - { // ask for next chunk CompletionStageDataActor.Chunk nextChunkCS AskPattern.ask( dataActor, (ActorRefDataActor.Chunk ref) - new DataActor.FetchChunk(currentOffset, ref), askTimeout, system.scheduler()); return nextChunkCS.thenApply( chunk - { ByteString bytes chunk.bytes; if (bytes.isEmpty()) return Optional.empty(); else return Optional.of(Pair.create(currentOffset bytes.size(), bytes)); }); });运行逻辑拆解初始状态为startOffset 0L第一次ask请求 offset 0 处的数据块收到Chunk后在map/thenApply中检查bytes.isEmpty非空返回Some((currentOffset bytes.length, bytes))把新状态推进到currentOffset bytes.length同时发射bytes下一次调用将用新 offset 继续拉取下一块为空返回None/Optional.empty()触发流的正常完成由于状态是Long不可变每次材质化都从startOffset重新开始行为完全确定。这样unfoldAsync就把带偏移量的 Actor 拉取协议平滑地折叠成了一个标准的、带背压的Source[ByteString, NotUsed]下游可以无缝衔接Sink、map、grouped等任意流操作符。更多使用场景与注意事项场景一异步斐波那契数列在源码注释与单元测试中都出现了用unfoldAsync生成斐波那契数列的示例。测试代码位于 SourceSpec.scalagenerate a finite fibonacci sequence asynchronously in { Source .unfoldAsync((0, 1)) { case (a, _) if a 10000000 Future.successful(None) case (a, b) Future(Some((b, a b) - a))(system.dispatcher) } .runFold(List.empty[Int]) { case (xs, x) x :: xs } .futureValue should (expected) }与unfold的同步版本同文件Source.unfold((0, 1))(...)的generate an unbounded fibonacci sequence用例以及 scaladsl/Source.scala 中成对的文档注释相比unfoldAsync的折叠函数在每次迭代都经过一个Future这正是它适合异步计算形态的体现。场景二无限数据源和unfold一样如果折叠函数永远不返回None/Optional.empty()unfoldAsync将产生无限流必须配合.take(n)等操作符截断。这在有界分页拉取场景中尤其重要——例如在循环分页时务必设计好数据耗尽的判定如返回空页否则流将永不完成。注意事项状态应不可变初始状态s会被复用于每次材质化若状态是可变对象如java.util.Iterator、Array、Java 标准集合多次材质化可能互相污染文档建议结合Source.lazySource让每次材质化创建全新的可变状态异步结果不能为 nullScala 的Future不能持有null的OptionJava 端CompletionStage的结果必须是合法对象Optional不能为null否则处理逻辑会失效Future 失败即流失败如内部实现所示Failure(ex)会直接fail(out, ex)因此应在折叠函数内部map/thenApply把业务上的边界情况映射为None/Optional.empty()而不是让 Future 失败每次调用只产生一个元素unfoldAsync每次异步结果最多发射一个元素如果需要一次异步请求产出多个元素应考虑mapConcat、flatMapConcat等操作符的组合阻塞 I/O 优先用 unfoldResource对于同步阻塞式资源读取官方文档建议使用unfoldResource基于同步阻塞调用而unfoldAsync面向异步 API。小结Source.unfoldAsync是连接异步状态推进与Reactive Streams 背压之间的桥梁签名简单(S) Future[Option[(S, E)]]/FunctionS, CompletionStageOptionalPairS,E无需接触低层GraphStage语义清晰Future 完成且结果为Some时发射并推进状态为None时完成失败时流失败且发射/完成均受下游需求驱动实现可靠底层UnfoldAsync/UnfoldAsyncJavaGraphStage 通过getAsyncCallback保证异步回调安全地回到流执行线程见 impl/Unfold.scala实战价值高从 Actor ask 分块读取、异步斐波那契到任意按游标/偏移量分页拉取的异步数据源它都能以声明式方式表达并天然获得 Akka Streams 的背压、取消与错误传播保障。若需要继续深挖可对照阅读同步版本 Source.unfold 文档 与对应的 Java 测试 SourceTest.java二者在状态机与语义上高度一致有助于完整理解整个 unfold 操作符家族。赞分享后端并发编程异步编程【免费下载链接】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 scanAsync 操作符详解基于 Future/CompletionStage 的异步累加扫描Akka Streams scanAsync 操作符详解基于 Future/CompletionStage 的异步累加扫描 scanAsync 是 Akka后端并发编程异步编程Akka Streams Sink.futureSink 详解将 Future[Sink] 接入流式数据消费Akka Streams Sink.futureSink 详解将 Future Sink 接入流式数据消费 导读 Sink.futureSink 是 Akka后端并发编程异步编程Akka Streams Flow.completionStageFlow基于 CompletionStage 的延迟 Flow 创建与流式接入指南Akka Streams Flow.completionStageFlow基于 CompletionStage 的延迟 Flow 创建与流式接入指南 本篇技术后端并发编程异步编程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考