ARTICLE DETAIL

资讯详情

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

Akka Streams 算子详解:Source.fromFutureSource 的弃用与 futureSource 迁移指南

Akka Streams 算子详解:Source.fromFutureSource 的弃用与 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点击查看免费下载导读Source.fromFutureSource是 Akka Streams 中用于等待一个 Future 完成后再将其内部的 Source 展开为流的工厂算子。它曾在需要异步获取数据源例如远程服务连接建立后再取得用户数据流的场景中扮演关键角色但在 Akka 2.6.0 起已被正式弃用官方推荐使用语义完全一致的 Source.futureSource 替代。本文以官方算子文档为主体结合仓库内 Scala/Java 双 API 的实现源码与测试用例梳理该算子的功能语义、弃用原因、迁移方式、底层实现原理与 Reactive Streams 背压语义帮助你在新代码中正确选用算子并平滑完成旧代码迁移。一、算子定位与弃用状态fromFutureSource属于 Source 算子集合Source operators位于akka.stream.scaladsl.Source伴生对象中。官方文档明确声明Deprecated bySource.futureSource.该算子在2.6.0版本中被标记为deprecated官方给出的迁移指引为改用futureSource而在 Java DSL 中则提示Use Source.futureSource (potentially together with Source.fromGraph)instead见 javadsl/Source.scala 第 210 行。弃用本身并不代表功能被移除或替换为不同行为——fromFutureSource与futureSource的核心语义完全一致仅仅是命名与签名的规范化调整因此迁移成本极低多数情况下只需将方法名从fromFutureSource改为futureSource。Scala 签名def fromFutureSourceT, M: Source[T, Future[M]]在 scaladsl/Source.scala 中其替换者futureSource的签名与之对等只是入参类型从Graph[SourceShape[T], M]收紧为更常用的Source[T, M]def futureSourceT, M: Source[T, Future[M]]两者的返回类型都是Source[T, Future[M]]即新 Source 的物化值materialized value是一个Future[M]它会在外层 Future 完成后携带内层 Source 的物化值 M。Java 签名Java DSL 中对应的弃用签名接受Future[_ : Graph[SourceShape[T], M]]见 javadsl/Source.scala其实现仅是对 Scala 版本的薄封装deprecated(Use Source.futureSource (potentially together with Source.fromGraph) instead, 2.6.0) public static T, M SourceT, FutureM fromFutureSource(Future? extends GraphSourceShapeT, M future)Java 使用者应迁移到Source.futureSource若入参是 Java 标准库的CompletionStage则应使用对应的 completionStageSource 算子详见下文。二、核心语义等待 Future随后展开内层 Source无论文档还是源码对fromFutureSource/futureSource语义的描述都高度一致Streams the elements of the given future source once it successfully completes. If the future fails the stream is failed.即传入一个Future[Source[T, M]]当该 Future成功完成后流开始发射内层 Source 中的元素如果该 Future失败整个流随之失败并把异常传递给下游。它解决的是典型的异步获取数据源问题数据源本身要等待某个异步操作如建立连接、鉴权、加载配置完成后才可用。官方示例见 FutureSource.scala给出的场景是通过 HTTP/2 或 WebSocket 访问远程服务把远程用户数据建模为Source[User, NotUsed]但该 Source 只有等连接建立后才可用import akka.NotUsed import akka.stream.scaladsl.Source import scala.concurrent.Future // 远程服务抽象连接建立后返回用户数据流 trait UserRepository { def loadUsers: Future[Source[User, NotUsed]] } // 等待 Future 完成随后展开为用户数据流 val userFutureSource: Source[User, Future[NotUsed]] Source.futureSource(userRepository.loadUsers)在这个例子中loadUsers返回的 Future 封装了连接建立这一异步步骤futureSource负责在连接就绪后无缝地把用户数据流接入主流水线。三、Reactive Streams 语义官方文档给出了明确的背压与完成语义callout 块语义描述emits一旦外层 Future 完成发射其内层future source的下一个值completes当内层future source完成时整个流完成需要补充说明的边界行为还包括失败传播外层 Future 以失败结束时流立即失败不会发射任何元素内层流特性透传内层 Source 的完成、失败以及背压信号都会原样传递给下游futureSource不会改变内层流的发射节奏物化值新流的物化值为Future[M]在下层 Source 物化后完成。四、源码级原理futureSource 的三条快路径与 FutureFlattenSource阅读 scaladsl/Source.scala 第 545-551 行 可以发现futureSource的实现针对 Future 的不同状态做了分流优化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.successful零额外开销Future 已失败直接构造一个Source.failed(exc)物化值为Future.failed(exc)同样避免建立任何中间机制Future 尚未完成进入通用路径构建FutureFlattenSource图阶段等待 Future 完成后动态衔接内层 Source。4.1 FutureFlattenSource 的展开机制FutureFlattenSource定义于 impl/fusing/GraphStages.scala 第 340 行起是一个GraphStageWithMaterializedValue[SourceShape[T], Future[M]]final class FutureFlattenSourceT, M extends GraphStageWithMaterializedValue[SourceShape[T], Future[M]] { val out: Outlet[T] Outlet(FutureFlattenSource.out) ... }其核心设计要点如下物化值即未来结果createLogicAndMaterializedValue中通过Promise[M]()生成物化 Future当内层 Source 物化时该 Promise 被完成preStart 快路径preStart()首先检查futureSource.value若 Future 已就绪则直接走同步回调注释说明这是为了avoid going through any execution context, in similar vein to FastFuture即规避线程池调度开销否则注册getAsyncCallback通过ExecutionContext.parasitic在 Future 完成时触发回调下游取消保护初始 OutHandler 的onDownstreamFinish中如果下游在 Future 完成前取消且物化 Promise 尚未完成则以StreamDetachedException(Stream cancelled before Source Future completed)失败该 Promise。源码注释特别指出早期实现曾在此时尝试物化内层 Source 以取得物化值但那不安全、可能导致 graph shell 泄漏因此被移除SubSinkInlet 衔接Future 完成后阶段通过SubSinkInlet[T]把内层 Source 作为子流接入onPush/onPull实现标准的推拉背压保证内层流的元素逐个透传到out。这段实现从图阶段GraphStage层面印证了文档语义在外层 Future 完成之前阶段不发射任何元素完成之后内层流的行为发射、完成、失败、背压被原样透传。五、测试用例对语义的验证仓库中 SourceSpec.scala 的Source.futureSource must套件覆盖了文档声明的全部关键行为可直接作为行为契约参考已完成的 Future 优化Future.successful(Source.single(done))时结果 Source 的算子名被优化为singleSource而不是futureSource证明走了快路径透传物化值透传对已完成 Future 的内层 Source 调用mapMaterializedValue(_ materializedValue)后外层流的物化值仍能取得materializedValue已失败的 FutureFuture.failed[Source[String, NotUsed]](TE(boom))时流的输出与物化 Future都以TE(boom)失败印证future fails the stream is failed延迟失败的 Future通过Promise在物化后再promise.failure(TE(boom))流输出与物化值同样双双失败不重复取消子流futureSource(akka.pattern.after(2.seconds)(...))与另一 Sourcemerge后take(1)验证下游提前取消时不会对子流产生双重取消。此外DslFactoriesConsistencySpec见 akka-stream-tests/src/test/scala/akka/stream/DslFactoriesConsistencySpec.scala会校验 Scala/Java DSL 工厂方法的一致性保证futureSource在两个 DSL 中行为对齐。六、迁移到 futureSource 与 completionStageSource6.1 Scala直接改名由于签名等价Scala 中的迁移几乎只是方法名替换// 迁移前2.6.0 起弃用 val source: Source[User, Future[NotUsed]] Source.fromFutureSource(userRepository.loadUsers) // 迁移后 val source: Source[User, Future[NotUsed]] Source.futureSource(userRepository.loadUsers)6.2 JavafutureSource 与 completionStageSourceJava DSL 中Source.futureSource接受CompletionStage或Future包装的 Source。若你持有的是 Java 标准库CompletionStage[Source[T, M]]应使用专门的completionStageSource算子见 javadsl/Source.scala 第 310-318 行其返回Source[T, CompletionStage[M]]实现上通过ExecutionContext.parasitic将CompletionStage转为 ScalaFuture后委托给futureSource再把物化值转回CompletionStage// Java使用 CompletionStage 版本 CompletionStageSourceUser, NotUsed futureSource repo.loadUsers(); SourceUser, CompletionStageNotUsed source Source.completionStageSource(futureSource);七、适用场景与相关算子对比futureSource及弃用的fromFutureSource典型适用于数据源本身是异步就绪的且就绪后以流的形态消费例如 WebSocket/HTTP 连接建立、按需初始化数据库游标、异步加载后展开分页流等。与易混淆算子区分算子输入行为物化值Source.futureFuture[T]Future 完成后发射单个元素并完成NotUsedSource.futureSource/fromFutureSourceFuture[Source[T, M]]Future 完成后展开为内层流Future[M]Source.lazySource() Source[T, M]有下游需求时才调用工厂创建 Source惰性Future[M]若需要延迟到有下游需求才创建而非立即开始等待 Future应选用lazySource系列见 scaladsl/Source.scala 中lazySource/lazyFutureSource的文档说明二者在何时触发这一维度上存在本质区别。八、小结Source.fromFutureSource自Akka 2.6.0起弃用语义由Source.futureSource完整承接两者签名对等、行为一致迁移只需替换方法名算子语义可概括为外层 Future 成功则展开内层流失败则整个流失败且内层流的完成、失败与背压信号原样透传实现上通过已完成/已失败快路径 FutureFlattenSource图阶段兼顾零开销优化与通用异步等待并在下游提前取消时以StreamDetachedException安全终止物化 PromiseJava 用户根据入参类型选择futureSource或completionStageSource相关行为均有源码与测试双重背书可在 scaladsl/Source.scala、GraphStages.scala 与 SourceSpec.scala 中进一步查阅。赞分享后端并发编程异步编程【免费下载链接】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 的 Source.fromSourceCompletionStage弃用原因与 completionStageSource 迁移指南Akka Streams 的 Source.fromSourceCompletionStage弃用原因与 completionStageSource 迁移指南后端并发编程异步编程Akka Streams Sink.lazyInitAsync 详解按首个元素延迟创建 Sink 的已弃用算子及其迁移方案Akka Streams Sink.lazyInitAsync 详解按首个元素延迟创建 Sink 的已弃用算子及其迁移方案 本指南围绕 Akka Stream后端并发编程异步编程Akka Streams FileIO.toFile 文件写入 Sink 详解已废弃 API 的用法、内部实现与 toPath 迁移指南Akka Streams FileIO.toFile 文件写入 Sink 详解已废弃 API 的用法、内部实现与 toPath 迁移指南 导读 FileIO.后端并发编程异步编程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表