ARTICLE DETAIL

资讯详情

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

Akka Streams Source.zipN 详解:将多个上游源合并为元素序列流

Akka Streams Source.zipN 详解:将多个上游源合并为元素序列流 Akka Streams Source.zipN 详解将多个上游源合并为元素序列流【免费下载链接】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.zipN是 Akka Streams 中用于多路合并fan-in的Source组合算子它接收任意数量的上游 Source按序从每个上游各取一个元素打包成一个序列Scala 为immutable.Seq/VectorJava 为List后向下游发射。本文以官方算子文档 zipN.md 为核心结合 Scala DSL 源码、GraphStage 实现与 官方测试样例 展开。读完本文你将掌握zipN的签名与行为语义、Scala/Java 双端调用方式、源码级运行原理含背压与完成时机以及与zip、zipWith、zipAll、zipWithN等兄弟算子的选型差异。核心语义什么是 Source.zipNSource.zipN将多个 Source 的元素按索引配对地组合成一个新的 Source其下游每个元素都是一个由各上游元素组成的序列。该算子属于 Source operators 索引 中的标准内建算子。它的行为可概括为三点每次发射要求所有上游各就绪一个元素——当且仅当全部输入端口都有元素可用时才把这一组元素按下游发射下游序列的元素顺序与传入的 sources 列表顺序完全一致任一上游结束整个zipN立即结束——表现为木桶效应最终结果长度取决于最短的上游。由于 sources 是以列表形式传入的各源的静态类型在列表中被抹平Scala 端下游序列会包含所有源元素的最近公共超类型closest supertypeJava 端则需要你自己把各源向上转型为共同的父类型后再调用zipN。签名Scala DSL位于 akka-stream/src/main/scala/akka/stream/scaladsl/Source.scaladef zipNT: Source[immutable.Seq[T], NotUsed]Java DSL位于 akka-stream/src/main/scala/akka/stream/javadsl/Source.scalapublic static T SourceListT, NotUsed zipN(ListSourceT, ? sources)从签名可以看到两点物化值被折叠为NotUsed传入的多个 Source 各自可能带有物化值如Source.queue的SourceQueue但zipN只返回组合后流的物化值NotUsed中间源的物化值无法再被访问元素类型被统一所有输入必须能视为同一类型T输出为immutable.Seq[T]Java 为List[T]。实战示例字符、数字与颜色的三路合并官方测试样例同时给出了 Scala 与 Java 两种写法Zip.scala、Zip.java。Scala 示例import akka.actor.typed.ActorSystem import akka.stream.scaladsl.Source implicit val system: ActorSystem[_] ??? val chars Source(a :: b :: c :: e :: f :: Nil) val numbers Source(1 :: 2 :: 3 :: 4 :: 5 :: 6 :: Nil) val colors Source(red :: green :: blue :: yellow :: purple :: Nil) Source.zipN(chars :: numbers :: colors :: Nil).runForeach(println) // prints: // Vector(a, 1, red) // Vector(b, 2, green) // Vector(c, 3, blue) // Vector(e, 4, yellow) // Vector(f, 5, purple)Java 示例import akka.NotUsed; import akka.actor.ActorSystem; import akka.stream.javadsl.Source; import java.util.Arrays; import java.util.List; ActorSystem system null; SourceObject, NotUsed chars Source.from(Arrays.asList(a, b, c, e, f)); SourceObject, NotUsed numbers Source.from(Arrays.asList(1, 2, 3, 4, 5, 6)); SourceObject, NotUsed colors Source.from(Arrays.asList(red, green, blue, yellow, purple)); Source.zipN(Arrays.asList(chars, numbers, colors)).runForeach(System.out::println, system); // prints: // [a, 1, red] // [b, 2, green] // [c, 3, blue] // [e, 4, yellow] // [f, 5, purple]注意 Java 示例中的细节三个源分别被声明为SourceObject, NotUsed这正是文档中提到的Java 端需要先将各源转型为共同超类型——chars与colors是字符串流、numbers是整型流它们的公共父类型是Object因此 Java 端必须显式统一类型后才能放入同一个List调用zipN。观察输出规律每个输出元素都是三元组且位置与传入顺序严格对应第一位永远来自chars第二位来自numbers第三位来自colorschars与colors各只有 5 个元素而numbers有 6 个。输出恰好 5 行——numbers的第 6 个元素6永远不会被消费因为当chars和colors发射完第 5 个元素后即完成zipN随之完成。这正是completes when any upstream completes语义的直观体现。源码级原理从 zipN 到 ZipWithN GraphStagezipN并非独立实现而是建立在更通用的zipWithN之上的特例。我们沿调用链逐层拆解所有行号均指向 akka-stream/src/main/scala/akka/stream/scaladsl/Source.scala。第一层zipN 委托给 zipWithN// L831-832 def zipNT: Source[immutable.Seq[T], NotUsed] zipWithN(ConstantFun.scalaIdentityFunction[immutable.Seq[T]])(sources).addAttributes(DefaultAttributes.zipN)zipN等价于以恒等函数seq seq作为 zipper 的zipWithN并额外附加DefaultAttributes.zipN对应 Stages.scala 中的命名属性用于调试与算子统计。第二层zipWithN 的三种分支// L837-846 def zipWithNT, O(sources: immutable.Seq[Source[T, _]]): Source[O, NotUsed] { val source sources match { case immutable.Seq() empty[O] case immutable.Seq(source) source.map(t zipper(immutable.Seq(t))).mapMaterializedValue(_ NotUsed) case s1 : s2 : ss combine(s1, s2, ss: _*)(ZipWithN(zipper)) case _ throw new IllegalArgumentException() // just to please compiler completeness check } source.addAttributes(DefaultAttributes.zipWithN) }从源码可以确认三个边界分支| 输入源数量 | 行为 | |--|--| | 0 个源 | 直接返回Source.empty流立即完成、零发射 | | 1 个源 | 用map把每个元素包成单元素序列物化值折叠为NotUsed| | ≥ 2 个源 | 通过combine将全部源接入ZipWithN这个GraphStage|第三层ZipN 是带恒等 zipper 的 GraphStageZipWithN与ZipN定义在 akka-stream/src/main/scala/akka/stream/scaladsl/Graph.scala// L1172-1175 final class ZipNA extends ZipWithN[A, immutable.Seq[A]](ConstantFun.scalaIdentityFunction)(n) { override def initialAttributes DefaultAttributes.zipN override def toString ZipN }ZipN是ZipWithN的恒等特例而ZipWithN是一个GraphStage[UniformFanInShape[A, O]]其形状shape为UniformFanInShapeA, O——即n 个同类型输入端口、1 个输出端口。Scala DSL 的combine会把传入的 n 个源逐一~连接到对应输入端口上。第四层GraphStageLogic 的运行机制ZipWithN.createLogic中的关键状态机逻辑Graph.scala L1203-L1247var pending 0 var willShutDown false ... override def preStart(): Unit shape.inlets.foreach(pullInlet) // 启动时向所有输入拉取 def onPull(): Unit { pending n; if (pending 0) pushAll() } // 下游每拉取一次登记 n 个待收元素 // 每个输入端口 override def onPush(): Unit { if (i 0) contextPropagation.suspendContext() pending - 1 if (pending 0) pushAll() // 收齐 n 个元素才发射 } override def onUpstreamFinish(): Unit { if (!isAvailable(in)) completeStage() willShutDown true // 任一上游完成即标记关闭 } private def pushAll(): Unit { contextPropagation.resumeContext() push(out, zipper(shape.inlets.map(grabInlet))) // 按输入端口顺序 grab 并打包 if (willShutDown) completeStage() else shape.inlets.foreach(pullInlet) // 发射后继续向所有输入拉取 }从中可以提炼出实现层面的结论栅栏barrier语义由pending计数器实现下游每产生一次需求onPull就登记n个待收元素只有 n 个输入端口全部onPush之后pending 0才调用pushAll发射。因此只要有一个上游慢其它已就绪的上游就会一直持有元素等待——这就是文档中backpressures 所有上游的底层来源发射顺序依赖shape.inlets.map(grabInlet)inlets按端口索引排列与传入sources的顺序一致从而保证输出序列的元素顺序与源列表顺序相同完成时机的微妙处理onUpstreamFinish中若当前没有正在等待被 grab 的元素!isAvailable(in)则直接completeStage()否则仅置willShutDown true待当前批次pushAll发射完这一组完整元素后再完成。注释说明这样可避免多一次多余的 pull保证已凑齐的整组元素仍会被完整发射然后立即结束上下文ContextPropagation传播从第一个输入端口i 0挂起上下文并在pushAll时恢复保证穿过该 stage 的上下文延续性。Reactive Streams 语义官方文档给出的契约可对照上文源码验证emits发射当所有输入端口都有元素可用时发射由各输入元素组成的序列completes完成当任意上游完成时完成即最短源决定流的总长度backpressures背压当下游背压时会背压所有上游同时某个上游即使已发射过元素也会一直被背压到其余所有上游都发射了各自的元素栅栏等待对应pending计数逻辑。边界场景与实用注意事项元素类型向上转型因为输入是列表zipN无法保留各源的精确元素类型。Scala 中下游元素类型是最小公共超类型Java 中必须先手动把源统一转型为公共父类型见上文 Java 示例的SourceObject, NotUsed。物化值丢失返回类型恒为Source[Seq[T], NotUsed]输入源自身的物化值不可达。若需要访问物化值请在调用zipN之前先物化各源或改用其它组合方式。最短源决定长度若各源长度不齐超出最短源长度的元素永远不会被消费也不会被拉取因此不会产生额外开销。适合的输入规模zipN面向多个源≥2的统一打包场景若只需合并两个源可直接用zip若需要在打包时做聚合转换应优先考虑zipWithN本文示例中zipWithN((seq: Seq[Int]) seq.max)即为取三者最大值的用法见 Zip.scala。与相关算子的对比选型zipN属于 zip 家族在文档的 See also 中列出了全部兄弟算子建议按需求选择| 算子 | 输入 | 输出 | 适用场景 | |--|--|--|--| | zipN | n 个源 |Seq[T]| 任意数量源按位合并成序列 | | zipWithN | n 个源 |O自定义 | 合并 n 个源并立即做聚合zipN即其恒等特例 | | zip | 2 个源 |(A, B)二元组 | 固定两个源的按位配对 | | zipAll | 2 个源 |(A, B)二元组 | 允许较短源结束后用默认值补位而非立即完成 | | zipWith | 2 个源 |O自定义 | 两个源按位合并并应用转换函数 | | zipWithIndex | 1 个源 |(T, Long)| 为元素附带递增序号 |关键取舍在于完成策略zipN/zipWithN/zip/zipWith都是任一上游完成即整体完成而zipAll允许通过默认值补齐继续发射当需要把 N 个源的同一批次聚合成一个结果如求最大、拼接、求和时zipWithN比zipN后再map更直接高效。小结Source.zipN以极简的 API 解决了多路源按位打包这一高频合并需求Scala/Java 双端签名统一、输出顺序与输入顺序严格一致、栅栏式背压保证数据对齐、最短源决定生命周期。从源码看它是通用zipWithN在恒等函数下的特例底层由ZipWithNGraphStage 的pending计数状态机驱动理解这一实现细节有助于在实际项目中准确预判它的完成时机、背压行为与类型约束。若读者想继续深入可阅读其实现源码 Source.scala 与 Graph.scala或运行 Zip.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创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表