后端并发编程异步编程【免费下载链接】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点击查看免费下载Sink.seq是 Akka Streams 中一个常用且易用的下游算子Sink operator它把上游流发出的全部元素逐个收集进一个集合并在流正常完成时通过物化值materialized value把集合返回给调用方。本文以 seq.md 文档为核心结合akka-stream模块中的源码实现与测试用例系统讲解Sink.seq的签名、用法、底层原理、边界行为与内存考量读完你可以直接在项目中使用它也能理解它为什么“能收集所有元素”以及什么时候会主动取消上游流。功能概述收集流中的所有元素Sink.seq会持续接收上游upstream发出的每一个元素将其追加到内部缓冲区直到上游流完成。当流正常完成时所有收集到的元素会作为一个整体集合通过物化值交还给调用方Scala 侧物化值为Future[Seq[T]]具体运行时通常为不可变的VectorJava 侧物化值为CompletionStageListT。也就是说你不需要自己维护一个可变的List再手动addSink.seq把“收集结果”这件事内置进了流图的物化过程让整条流的处理保持声明式风格。需要特别注意的一点官方文档与源码 scaladoc 都明确强调集合的大小受限于Int.MaxValueJava 为Integer.MAX_VALUE。如果上游发出的元素超过了这个上限Sink.seq会取消cancel整个流而不是无限增长内存。签名与物化值类型Scala 签名在 akka-stream/src/main/scala/akka/stream/scaladsl/Sink.scala 中定义def seq[T]: Sink[T, Future[immutable.Seq[T]]]它不接收任何参数类型参数T由上游元素类型推断。物化值是Future[immutable.Seq[T]]该Future在流完成时以收集到的集合成功完成若流失败则Future以对应异常失败。Java 签名在 akka-stream/src/main/scala/akka/stream/javadsl/Sink.scala 中定义def seq[In]: Sink[In, CompletionStage[java.util.List[In]]]Java API 的物化值是CompletionStageListIn内部通过scaladsl.Sink.seq包装并借助CollectionConverters把 Scala 的Seq转换为java.util.List转换过程使用ExecutionContext.parasitic不引入额外的异步调度开销。物化值的语义物化Future/CompletionStage只有三种结局流正常完成→ 以收集到的全部元素成功完成流失败上游发出错误→ 以该异常失败流被异常中止如因取消、急停导致 stage 提前停止→ 以AbruptStageTerminationException失败见下文源码。实战示例收集数字流官方文档分别给出了 Scala 与 Java 两个可直接运行的示例。Scala 示例摘自 SinkSpec.scala 中The seq sink must测试块val source Source(1 to 3) val result source.runWith(Sink.seq[Int]) val seq result.futureValue seq.foreach(println) // will print // 1 // 2 // 3 assert(seq Vector(1, 2, 3))Source(1 to 3)依次发出1、2、3三个元素Sink.seq[Int]将它们全部收集流完成后result一个Future[Seq[Int]]成功完成其值为Vector(1, 2, 3)。注意 Scala 侧默认收集到的具体集合类型是不可变Vector对应源码中的new SeqStage[T, Vector[T]]。Java 示例摘自 SinkDocExamples.javaSourceInteger, NotUsed ints Source.from(Arrays.asList(1, 2, 3)); CompletionStageListInteger result ints.runWith(Sink.seq(), system); result.thenAccept(list - list.forEach(System.out::println)); // 1 // 2 // 3Java 侧通过result.thenAccept(...)异步消费物化出来的ListInteger同样打印1、2、3。典型使用场景批处理结果收集把一批待处理数据流入流式管线处理后统一收集结果再做聚合或落库测试断言如SinkSpec所示配合futureValue将流结果取回后与期望集合直接比对是 Akka Streams 测试中最常见的断言方式之一流的“汇点”作为Source.runWith或flow.runWith(Sink.seq)的终点把流式处理收敛为一次性的异步结果。底层实现SeqStage 源码剖析Sink.seq的实现非常轻量本质上是一个自定义的GraphStageWithMaterializedValue。理解它的实现有助于你精确把握“何时成功、何时失败、何时取消”的语义。工厂方法与默认属性在 Sink.scala 中def seq[T]: Sink[T, Future[immutable.Seq[T]]] Sink.fromGraph(new SeqStage[T, Vector[T]])同时还有一个更通用的变体Sink.collection[T, That]它通过隐式Factory允许你指定目标集合类型例如Seq、Vector等。该 stage 的默认属性在 Stages.scala 中注册为seqSink名称用于调试与日志输出。SeqStage 的核心逻辑完整实现在 akka-stream/src/main/scala/akka/stream/impl/Sinks.scala关键点如下InternalApi private[akka] final class SeqStageT, That extends GraphStageWithMaterializedValue[SinkShape[T], Future[That]] { val in InletT // ... override def createLogicAndMaterializedValue(...) { val p: Promise[That] Promise() val logic new GraphStageLogic(shape) with InHandler { val buf cbf.newBuilder override def preStart(): Unit pull(in) def onPush(): Unit { buf grab(in) pull(in) } override def onUpstreamFinish(): Unit { val result buf.result() p.trySuccess(result) completeStage() } override def onUpstreamFailure(ex: Throwable): Unit { p.tryFailure(ex) failStage(ex) } override def postStop(): Unit { if (!p.isCompleted) p.failure(new AbruptStageTerminationException(this)) } setHandler(in, this) } (logic, p.future) } }这段代码揭示了完整的数据流与生命周期语义预启动拉取preStart()中立即pull(in)stage 一启动就向上游请求第一个元素不浪费任何吞吐机会逐元素追加每次onPush()把grab(in)到的元素追加进cbf.newBuilder构建的缓冲区然后继续pull(in)请求下一个元素——这是典型的“推-拉”回压循环上游按下游的拉取节奏逐步供数天然具备背压backpressure能力正常完成onUpstreamFinish()中buf.result()冻结出不可变集合p.trySuccess(result)完成物化Future再completeStage()结束整个 stage失败传播onUpstreamFailure(ex)把异常同时传给物化Future与流图failStage保证“流失败 ⇒ 结果 Future 失败”的一致性异常中止兜底postStop()中若Future尚未完成例如 stage 因取消或图急停被终止则以AbruptStageTerminationException失败避免调用方永远等待一个永不完成的Future。可以推断Int.MaxValue的上限来自 Scala 集合按Int索引的设计约束Vector/Seq的规模上限而非SeqStage单独设置的门槛Java 侧的Integer.MAX_VALUE同理对应java.util.List的容量上限。边界行为与取消语义Reactive Streams semantics官方文档在 “Reactive Streams semantics” 一节中给出的语义只有一条cancelsIf too many values are collected翻译过来即如果收集到的元素过多则取消上游流。结合上面的源码实现可以归纳出Sink.seq在 Reactive Streams 协议下的完整行为场景行为上游正常完成物化Future/CompletionStage成功完成携带全部收集元素上游失败物化结果以该异常失败stage 失败并向下游传播错误元素数量达到Int.MaxValue/Integer.MAX_VALUE取消上游流cancels阻止更多元素进入内存stage 被异常终止如取消/急停物化结果以AbruptStageTerminationException失败“取消”意味着Sink.seq不会在达到容量上限后继续吞入元素而是主动切断与上游的契约防止内存被无限耗尽——这是它作为“无界收集器”在极端情况下的安全阀。有界性考量用 take / limit 保护内存源码 scaladoc 对Sink.seq有一个重要提醒As upstream may be unbounded,Flow[T].takeor the stricterFlow[T].limit(and their variants) may be used to ensure boundedness.即上游可能是无界的。Sink.seq会一直收集到流结束如果上游是无限流如Source.repeat、Source.tick、Source.fromIterator配无限迭代器内存会持续增长。因此官方建议在Sink.seq之前显式限流例如// Scala只收集前 1000 个元素 Source(1L to Long.MaxValue) // 潜在无限/超大流 .take(1000) // 截断 .runWith(Sink.seq) // 收集前 1000 个// Java更严格的有界收集 ints.take(1000).runWith(Sink.seq(), system);官方文档与源码共同提到的相关有界化算子包括Flow.take取前 n 个元素后完成Flow.limit更严格的版本超过 n 个元素时以失败终止而不是静默截断Flow.limitWeighted按权重计数的限流Flow.takeWhile/Flow.takeWithin按谓词或时间窗口截断。在选择时take适合“截断即可”的场景limit适合“超过即报错、防止数据异常”的场景。与其他收集类 Sink 的对比Sink.seq属于“收集全部”型算子与akka.stream.scaladsl.Sink中的其他物化算子容易混淆从源码结构看它们各自定位不同算子物化值行为Sink.seqFuture[Seq[T]]收集全部元素流完成后返回完整集合Sink.headFuture[T]只取第一个元素流为空则失败Sink.headOptionFuture[Option[T]]取第一个元素流为空则返回NoneSink.lastFuture[T]只取最后一个元素Sink.takeLast(n)Future[Seq[T]]只保留最后 n 个元素Sink.foldFuture[S]需要初始值对元素做累积归约不保留中间元素Sink.collectionFuture[That]seq的泛化版本可指定目标集合类型选择原则很简单需要全量结果用seq只需要头/尾或聚合值时优先用head/last/fold等它们在内存占用上远优于seq因为不缓存全部元素。测试验证与进一步阅读Sink.seq的行为在仓库中有完整的测试覆盖Scala 测试akka-stream-tests/src/test/scala/akka/stream/scaladsl/SinkSpec.scala 验证了“收集流中的元素为序列”这一核心行为并断言结果等于Vector(1, 2, 3)Java 文档示例akka-docs/src/test/java/jdocs/stream/operators/SinkDocExamples.java 展示了 Java API 的完整用法。如果想要进一步探究算子的完整清单与分类可查看 stream operators 索引Sink.seq的 Scala 定义与相邻算子可查看 Sink.scala底层SeqStage的完整实现可查看 impl/Sinks.scala。小结Sink.seq是 Akka Streams 中“把流收敛为集合”的标准答案它声明式地收集全部元素通过Future/CompletionStage异步交付结果天然支持背压并在元素数量触及Int.MaxValue上限时主动取消上游以保证安全。使用时只需记住两个要点一是通过take/limit等算子显式约束上游规模以避免内存风险二是区分它与head/last/fold等只保留部分信息的算子——按需选择你的流式代码会更清晰、更健壮。赞分享后端并发编程异步编程【免费下载链接】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 Sink.collection 详解将流元素收集为任意 Scala 集合Akka Streams Sink.collection 详解将流元素收集为任意 Scala 集合 导读 Sink.collection 是 Akka Str后端并发编程异步编程Akka Streams Sink.takeLast 详解收集流末尾 n 个元素的实用指南Akka Streams Sink.takeLast 详解收集流末尾 n 个元素的实用指南 本指南以 Akka 官方文档 Sink.takeLast http后端并发编程异步编程Akka Streams zipWithIndex 算子详解为流元素自动编号的原理与实战Akka Streams zipWithIndex 算子详解为流元素自动编号的原理与实战 导读 zipWithIndex 是 Akka Streams 中一个后端并发编程异步编程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
