尧图网络 高端网站定制 · 原创设计
免费咨询热线
400-888-6620
免费获取方案
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),仅供参考
RELATED

相关推荐

Ceph OSD 与 Placement Group(PG)监控与故障排查实战指南

Ceph OSD 与 Placement Group(PG)监控与故障排查实战指南

存储分布式文件系统对象存储后端高可用 【免费下载链接】ceph Ceph is a distributed object, block, and file storage platform 项目地址: https://gitcode.com/gh_mirrors/ce/ceph 点击查看 免费下载 Ceph 是一个分布式对象、块与文件存储平台,其高…

📅 2026/9/23 18:48:17
论文选题避坑指南:3个实战项目教你搞定性能优化

论文选题避坑指南:3个实战项目教你搞定性能优化

论文选题避坑指南:3个实战项目教你搞定性能优化 官方文档翻了三遍,核心逻辑还是云里雾里?别急,这是大多数开发者的通病。 MDN Web Docs 里的示例代码往往过于理想化,直接复制到你的工程里,性能直接崩盘。 今天不聊虚的,直接上…

📅 2026/9/23 18:48:17
metaRTC8.0与metaIPC3.0:新架构下的IPC低延迟实践

metaRTC8.0与metaIPC3.0:新架构下的IPC低延迟实践

metaRTC到8.0这版,改动确实不小。我拿到metaIPC3.0之后,花了一周时间把原来测试环境里的几路摄像头全部切到新架构上跑了一遍,说实话,刚开始看代码结构的时候有点不适应,但理清楚之后你会发现,这套架构调整…

📅 2026/9/23 18:48:17
MORE NEWS

更多资讯

📰

广州体育生文化课冲刺学校哪家好?低分稳过线推荐

针对广州体育生长期专注术科训练、文化课学习断层、基础普遍偏弱、联考后复习时间紧张的备考现状,结合本地机构办学合规性、艺体生专项教学适配度、历年学员提分数据、市场真实口碑与精细化管理体系,适配体育生文化课冲刺的适配度较好的适配学校共有五家…

📰

Atlas 300V 24G部署YOLO完全指南:从推理加速卡选型到OM模型转换与性能调优

最近好几个做边缘视觉的朋友不约而同地问我同一个问题:Atlas 300V 24G到底是不是运算加速卡?能不能跑YOLO?我本来以为这是个随便搜搜就有答案的问题,但聊下来发现很多人都卡在“知道它叫Atlas,但不知道它和GPU有什么本…

📰

停车场空位检测数据集:VOC+YOLO双格式7959张工业级标注

简介:本资源为面向计算机视觉初学者与算法工程师的停车场空位检测专用数据集,适用于目标检测模型训练与评估,尤其适配YOLO系列及Pascal VOC兼容框架。数据集包含7959张高质量停车场实景图像,全部标注为“empty”和“occupied”两类…

📰

基于Python+OpenCV+Django的人脸识别课设源码全解析

简介:面向计算机相关专业课程设计与毕业设计场景,这份基于Python、OpenCV与Django的人脸识别系统源码,提供了一套从人脸检测、特征提取到浏览器端识别展示的完整落地方案。项目已通过导师指导并获得九十七分高分评价,代码完整且可…

📰

信用卡高风险识别毕业设计:从特征工程到模型实战

简介:这是一份基于Python实现的信用卡客户高风险识别毕业设计资源,面向计算机、人工智能、自动化等专业的在校学生、教师及企业员工,特别适合毕业设计、课程设计或实训作业场景。项目围绕历史信用风险、经济风险、收入风险三个维度构建客户属…

📰

基于OpenCV的车牌识别课程作业:HSV定位到字符分割与模板匹配

简介:一份基于Python3和OpenCV的数字图像处理课程作业车牌识别项目资料包,面向需要完成课程设计、大作业或入门图像处理的学习者,既可作为毕设/工程实训的初始框架,也适合有基础者二次改造。资源共19个文件,以Python源…

TODAY

今日更新

THIS WEEK

本周精选

THIS MONTH

本月热门

读完文章,想聊聊您的网站?

告诉我们您的行业与需求,资深顾问一对一梳理方案与报价,全程免费。

📞 💬