ARTICLE DETAIL

资讯详情

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

Akka Streams Source.completionStageSource 详解:等待异步就绪后再流入元素

Akka Streams Source.completionStageSource 详解:等待异步就绪后再流入元素 后端并发编程异步编程【免费下载链接】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.completionStageSource是 Akka Streams 中用于异步源就绪后再开始流动的关键算子它接收一个CompletionStageSource只有当这个异步阶段成功完成后才会把内部 Source 的元素转发给下游。它非常适合 HTTP/2、WebSocket、数据库连接池等连接建立后才能拿到数据流的场景本文将从官方文档、Java/Scala API 实现、底层 GraphStage 原理与测试验证四个层面完整讲解该算子的用法与机制帮助读者在真实项目中正确选用它。一、算子定位与适用场景Source.completionStageSource的核心语义是Streams the elements of an asynchronous source once its givencompletionoperator completes.即只有当传入的CompletionStage完成成功之后才会开始流式输出其内部异步源的元素。如果该CompletionStage以失败结束则整个流也会以该异常失败。典型场景访问一个通过 HTTP/2 或 WebSocket 提供用户数据流User data stream的远程服务。我们可以把远程数据流建模为Source[User, NotUsed]但这个 Source 只有在连接建立之后才真正可用。此时Source.completionStageSource就是连接建立异步完成与数据流动之间的桥梁。// 远程服务抽象loadUsers() 异步返回一个数据流 interface UserRepository { CompletionStageSourceUser, NotUsed loadUsers(); }对应的 Scala 侧标准库Future算子为Source.futureSource两者语义一致仅异步类型不同CompletionStage与Future。二、方法签名Java API 签名定义于 akka-stream/src/main/scala/akka/stream/javadsl/Source.scalastatic T, M SourceT, CompletionStageM completionStageSource( CompletionStageSourceT, M completionStageSource)参数与返回值要点项目说明输入CompletionStageSourceT, M异步完成的内部 Source输出SourceT, CompletionStageM元素类型与内部源一致物化值CompletionStageM内部 Source 物化完成后得到的物化值Scala 对应Source.futureSource三、官方示例远程用户数据流文档给出的完整 Java 示例位于 CompletionStageSource.javaimport akka.NotUsed; import akka.stream.javadsl.Source; import java.util.concurrent.CompletionStage; public class CompletionStageSource { public static void sourceCompletionStageSource() { UserRepository userRepository null; // an abstraction over the remote service SourceUser, CompletionStageNotUsed userCompletionStageSource Source.completionStageSource(userRepository.loadUsers()); // ... } interface UserRepository { CompletionStageSourceUser, NotUsed loadUsers(); } static class User {} }注意此处userRepository.loadUsers()的类型为CompletionStageSourceUser, NotUsed而最终得到的SourceUser, CompletionStageNotUsed——物化值从内部的NotUsed变为CompletionStageNotUsed这正是内部源物化时机异步化带来的类型变化。四、Reactive Streams 语义官方定义文档明确了该算子的背压语义emits发射当内部的异步源所关联的completion operator即传入的CompletionStage完成之后发射来自该异步源的下一个值completes完成当异步源完成时整个流完成。这意味着在CompletionStage完成之前下游的拉取pull/demand会被挂起一旦完成内部 Source 接入元素按 Reactive Streams 背压机制正常流动。五、源码级实现原理5.1 Java API 是对 Scala futureSource 的薄封装从 javadsl/Source.scala 可以看到completionStageSource的实现极为简洁def completionStageSourceT, M: Source[T, CompletionStage[M]] scaladsl.Source .futureSource(completionStageSource.asScala.map(_.asScala)(ExecutionContext.parasitic)) .mapMaterializedValue(_.asJava) .asJava关键点CompletionStage通过.asScala转换为 ScalaFuture内部javadsl.Source通过.asScala转换为scaladsl.Source转换动作运行在ExecutionContext.parasitic寄生执行上下文上——它不切换线程、直接在当前调用线程上执行避免额外调度开销最终物化值经.mapMaterializedValue(_.asJava)还原为CompletionStageM。5.2 Scala futureSource已完成 Future 的快速路径优化scaladsl/Source.scala 的实现包含快速路径fast path优化def futureSourceT, M: Source[T, Future[M]] { futureSource.value match { case Some(Success(source)) source.mapMaterializedValue(Future.successful) case Some(Failure(exc)) failed(exc).mapMaterializedValue(_ Future.failed(exc)) case _ fromGraph(new FutureFlattenSource(futureSource)) } }Future 已成功完成直接返回内部 Source并将物化值包装为已完成的Future完全跳过异步等待Future 已失败直接构造Source.failed(exc)流立即以该异常失败Future 尚未完成进入通用路径通过FutureFlattenSource图阶段GraphStage挂起等待。对应测试 SourceSpec.scala 验证了这些行为optimize already completed future in { val future Future.successful(Source.single(done)) val source Source.futureSource(future) source.getAttributes.nameLifted should (Some(singleSource)) // ... } handle already failed future in { val future Future.failed[Source[String, NotUsed]](TE(boom)) val source Source.futureSource(future) val (futureMat, streamResult) source.toMat(Sink.head)(Keep.both).run() streamResult.failed.futureValue should (TE(boom)) futureMat.failed.futureValue should (TE(boom)) }注意测试中nameLifted的断言当传入的 Future 已成功完成时算子的属性名直接变为内部源的名称如singleSource印证了快速路径确实直接替换为内部源而非新建包装节点。5.3 通用路径FutureFlattenSource GraphStage当异步阶段尚未完成时Akka 使用 FutureFlattenSourceGraphStageWithMaterializedValue[SourceShape[T], Future[M]]实现扁平化等待。其核心机制preStart()时再次检查futureSource.value——若已就绪则走类似快速路径的优化注释明确说明这是避免经过任何执行上下文与 FastFuture 同思路的优化否则通过getAsyncCallback[Try[Graph[SourceShape[T], M]]]注册异步回调并用ExecutionContext.parasitic订阅 Future 的完成回调内部使用SubSinkInlet[T]作为子源的入口——它本质上把一个 Source 当作 Sink 的反向接入从而复用完整的 Reactive Streams 订阅机制实现背压传递若 Future 失败则sinkIn.cancel()、物化 Promise 失败、failStage(t)使整个流失败物化值通过内部Promise[M]桥接当子源物化后完成。这段实现是理解Source.completionStageSource与flatMapConcat有何不同的关键它不做流的嵌套展开而是等待外部异步完成、再将单一内部源无缝接入并保证取消/失败等信号不会重复发送测试not cancel substream twice专门验证了这一点。六、与相关算子的对比与组合算子输入语义说明Source.completionStageSourceCompletionStageSourceT, M等待异步源就绪后流入元素本文主角Java APISource.futureSourceFuture[Source[T, M]]同上ScalaFuture版本futureSource.mdSource.completionStageCompletionStageT单值异步化仅流出一个元素见 javadsl/Source.scalaSource.lazyCompletionStageSourceCreatorCompletionStageSource延迟到有下游需求时才调用 create 获取异步源基于completionStageSource实现见 javadsl/Source.scalaSource.fromSourceCompletionStageCompletionStageGraph[SourceShape, M]旧 API接收 Graph自 2.6.0 起已deprecated建议改用completionStageSource必要时配合Source.fromGraph6.1 lazyCompletionStageSource更进一步如果创建异步源的代价很高例如每次建立 WebSocket 连接且希望直到有下游消费者时才真正发起连接可选用Source.lazyCompletionStageSourceSource.lazyCompletionStageSource( () - userRepository.loadUsers()) // 有需求时才调用其实现javadsl/Source.scala本质是lazySource[T, CompletionStage[M]](() completionStageSource(create.create())) .mapMaterializedValue(_.thenCompose(csm csm))即用lazySource包裹completionStageSource并把两层物化值通过thenCompose拍平。若下游在函数被调用前就取消或失败物化值将以NeverMaterializedException失败。七、实战要点与注意事项失败传播无论CompletionStage已失败还是稍后失败异常都会同时体现在两个层面——流的运行结果以及物化出的CompletionStageM见上文handle later failed future测试对futureMat与streamResult的双重断言。下游提前取消若在下游在内部源接入前取消物化出的CompletionStage将以StreamDetachedException失败javadsl 文档注释明确说明。避免不必要的异步等待completionStageSource内部对已完成的异步阶段有快速路径优化直接透传内部源不引入额外图节点——因此将已就绪的源包装进去几乎没有开销。与Source.from/Source.single组合如需等待的其实是一个Graph[SourceShape, M]按弃用提示应写成Source.completionStageSource(completion.thenApply(Source::fromGraph))这正是旧 APIfromSourceCompletionStage的内部实现见 javadsl/Source.scala。线程模型CompletionStage到Future的转换及内部回调均使用ExecutionContext.parasitic不引入线程切换符合 Akka Streams 对执行效率的一贯要求。八、总结Source.completionStageSource用最简洁的 API 解决了异步资源就绪后再启动数据流这一高频问题官方文档给出了清晰语义就绪前挂起、失败即失败、内部源完成后流完成源码实现则揭示了其背后快速路径 FutureFlattenSource 图阶段的高效机制而测试用例完整覆盖了已成功、已失败、稍后失败、物化值传递与取消去重等关键场景。在 Java 侧处理 WebSocket/HTTP/连接池等异步数据源时它是futureSource之外最值得优先使用的算子。赞分享后端并发编程异步编程【免费下载链接】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点击查看免费下载相关推荐Puppeteer WebWorker.waitForFunction() 方法详解在 Web Worker 运行时中等待异步条件就绪Puppeteer WebWorker.waitForFunction 方法详解在 Web Worker 运行时中等待异步条件就绪 本方法属于 WebWork浏览器控制测试网页爬虫开发工具Golongpoll安全最佳实践实现基于HTTP头的身份验证机制Golongpoll安全最佳实践实现基于HTTP头的身份验证机制 在当今的Web开发中实时通信已成为许多应用的核心需求。Golongpoll作为一款强大的G后端并发编程异步编程rn-sliding-up-panel与键盘交互完全指南实现完美用户体验的5种解决方案rn sliding up panel与键盘交互完全指南实现完美用户体验的5种解决方案 rn sliding up panel是一个基于React Nativ后端并发编程异步编程上一篇Pancake 开源项目教程下一篇React Native Raw Bottom Sheet 使用教程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表