首页/新闻资讯/正文详情

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

发布时间:2026/9/23 12:10:12 来源:云帆数科 栏目:资讯中心
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),仅供参考

相关推荐

MAX13487EESA自动方向RS-485收发器设计实战指南
MAX13487EESA自动方向RS-485收发器设计实战指南

简介:这是 MAX13487EESA 对应芯片 MAX13487E/MAX13488E 的英文原版数据手册,面向硬件工程师、嵌入式开发及工业通信调试人员,重点解决半双工 RS-485/RS-422 收发器选型与设计中的抗干扰、热插拔和自动方向控制问题。文件为 1 个 PDF&#xff… · 2026/9/23 12:10:12

百度文库如何免费下载保姆级教程:3个致命坑全解析
百度文库如何免费下载保姆级教程:3个致命坑全解析

百度文库如何免费下载保姆级教程:3个致命坑全解析 版本升级后 API 全变了,昨天还跑得通的脚本今天直接报 403 Forbidden,这种绝望感做过爬虫的都知道。别急着骂平台反爬严,90% 的情况是你没看懂它新的鉴权逻辑。这篇… · 2026/9/23 12:10:12

Servlet+JSP水电费管理系统:老旧小区可部署的Java毕业设计实战
Servlet+JSP水电费管理系统:老旧小区可部署的Java毕业设计实战

简介:本资源是一套面向计算机专业本科生的Java毕业设计实战项目——小区水电费管理系统,适用于课程设计、毕设选题与Java Web技术入门实践。系统基于B/S架构,采用JavaJSP开发,涵盖用户数据接入、水电费登记、在线统计与缴费管理四… · 2026/9/23 12:10:12

了不起的盖茨比英文图解原理:3步搞定环境配置卡死
了不起的盖茨比英文图解原理:3步搞定环境配置卡死

了不起的盖茨比英文图解原理:3步搞定环境配置卡死 配置环境就卡半天?别急,这通常是依赖地狱在搞鬼。 很多开发者对着终端报错发呆,其实核心逻辑没看透。 今天用图解原理拆解《了不起的盖茨比英文》数据处理的底层链路。… · 2026/9/23 12:56:46

JUnit 5扩展机制详解与实战应用
JUnit 5扩展机制详解与实战应用

1. JUnit 5扩展机制概述JUnit 5作为Java生态中最主流的测试框架,其扩展机制(Extension Model)是区别于旧版本的核心特性之一。不同于JUnit 4中通过Rule和Runner实现的有限扩展能力,JUnit 5通过统一的Extension API提供了更灵活的测… · 2026/9/23 12:56:45

猫咪ios下载避坑指南:API变更后的完整示例与薪资真相
猫咪ios下载避坑指南:API变更后的完整示例与薪资真相

猫咪ios下载避坑指南:API变更后的完整示例与薪资真相 版本升级后 API 全变了,很多转岗到 iOS 开发的兄弟直接懵圈。别慌,这里给你一份【猫咪ios下载】场景下的 完整示例 ,手把手拆解底层逻辑。 入口定位:从网络请求看转岗成本… · 2026/9/23 12:56:45

航拍牛羊小目标检测实战:从VOC/YOLO数据集到YOLOv8训练
航拍牛羊小目标检测实战:从VOC/YOLO数据集到YOLOv8训练

简介:面向智慧牧场航拍场景的牛羊检测任务,这份数据集聚焦远距离小目标识别难题,适合计算机视觉学习者、算法工程师及农业智能化项目开发者用于模型训练与算法验证。数据集中包含1021张航拍图像,每张均配套Pascal VOC格式的xml标注… · 2026/9/23 12:56:39

基于PyTorch的高分遥感语义分割实战:从数据准备到地物分类
基于PyTorch的高分遥感语义分割实战:从数据准备到地物分类

简介:面向遥感影像智能解译与深度学习语义分割开发者,这份PyTorch实现的“高分遥感语义分割(地物分类)”项目实践包,完整覆盖从数据准备、模型训练到预测后处理的流程。包内共858个文件,819张PNG图像承载了… · 2026/9/23 12:56:39

家书家训实战项目面试突击:3个高频考点拆解
家书家训实战项目面试突击:3个高频考点拆解

家书家训实战项目面试突击:3个高频考点拆解 很多开发者背熟了语法,却在面试实战项目中卡壳。 不是代码写不出,是逻辑理不清,痛点抓不准。 今天把【家书家训】相关高频题拆透,直击项目落地难点。 考点梳理:别只背定义,要看业务场景… · 2026/9/23 12:56:39

3招搞定手机怎么下载微信面试难题实战项目解析
3招搞定手机怎么下载微信面试难题实战项目解析

3招搞定手机怎么下载微信面试难题实战项目解析 面试被问“手机怎么下载微信”背后的原理,90%的人答不上来。别笑,这看似弱智的问题,实则是考察你对移动应用分发机制、安全校验及网络协议理解的试金石。我带过不少校招新人,他们背了八股文,却连一个A… · 2026/9/23 0:00:03

你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型
你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型

你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型 面试被问“高并发下如何保证消息不丢失”,你张口就是“用Redis”,结果面试官追问“如果Redis宕机了怎么办”,你瞬间卡壳。这种场景太常见了,很多新手在背八股文时,只记住了技术名词… · 2026/9/23 0:00:29

Win7无线热点配置工具源码解析:解决API失效的3个实战技巧
Win7无线热点配置工具源码解析:解决API失效的3个实战技巧

Win7无线热点配置工具源码解析:解决API失效的3个实战技巧 Win7无线热点配置工具在Win10/11上跑不动?不是你的问题,是版本升级后 API 全变了。很多老项目里的 netsh wlan… · 2026/9/23 0:00:36

了解更多?预约专属演示

我们的顾问将为您一对一讲解产品与方案

企业微信二维码