ARTICLE DETAIL

资讯详情

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

Akka Streams 的 alsoTo 算子:将元素旁路复制到附加 Sink 的 Fan-out 指南

Akka Streams 的 alsoTo 算子:将元素旁路复制到附加 Sink 的 Fan-out 指南 Akka Streams 的 alsoTo 算子将元素旁路复制到附加 Sink 的 Fan-out 指南【免费下载链接】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导读alsoTo是 Akka Streams 中一个轻量而实用的 Fan-out扇出算子它在把元素继续向下游传递的同时将同一份元素复制一份发送给一个附加的Sink。本指南以官方文档alsoTo.md为骨架结合akka-stream模块的源码实现与测试用例深入讲解它的签名、Reactive Streams 语义、与wireTap的区别、alsoToAll/alsoToMat变体以及典型实战场景。读完本文你将掌握如何在日志审计、指标采集、事件归档等场景中安全地旁路分流数据流并理解其背压行为对吞吐的影响。一、alsoTo 是什么一次传递两处送达alsoTo的核心语义在文档开头就给出了精确定义Attaches the givenSinkto thisFlow, meaning that elements that pass through thisFlowwill also be sent to theSink.即在原有数据流Source或Flow上挂接一个额外的Sink流经的元素会原样继续向下游传递同时也会被发送到该附加Sink。它属于 Fan-out operators 家族与broadcast、branch、divertTo等算子同属一路输入、多路输出的图形结构。在官方的算子分类索引中alsoTo与alsoToAll、divertTo、wireTap等一起被归入 Fan-out 类别适合在不改变主线数据流的前提下增加观察者式处理路径。底层实现一个两输出的 Broadcast从源码可以看到alsoTo并不是什么特殊魔法它的实现就是标准的Broadcast图拼接。在 scaladsl/Flow.scala 中def alsoTo(that: Graph[SinkShape[Out], _]): Repr[Out] via(alsoToGraph(that)) protected def alsoToGraphM: Graph[FlowShape[Out uncheckedVariance, Out], M] GraphDSL.createGraph(that) { implicit b r import GraphDSL.Implicits._ val bcast b.add(BroadcastOut) bcast.out(1) ~ r FlowShape(bcast.in, bcast.out(0)) }这段代码揭示了几个关键实现事实算子内部创建一个输出数为 2 的Broadcast端口0通向主下游端口1通向附加的SinkBroadcast使用了eagerCancel true只要任一输出取消广播即整体取消从而保证附加Sink与主下游的生命周期严格同步由于是Broadcast元素不是复制给两条路径各一份副本对象而是同一元素被分发到两个输出两个分支共享该元素整个拼接结果是一个FlowShape(in, out)通过via嵌入当前流因此alsoTo不会改变流的输入输出类型——元素类型仍为Out。二、签名与可用位置文档给出了 Scala 与 Java 两个 DSL 的签名Scaladef alsoTo(that: Graph[SinkShape[Out], _]): FlowOps.this.Repr[Out]Javadef alsoTo(that: Graph[SinkShape[Out], _]): javadsl.Flow[In, Out, Mat]alsoTo定义在FlowOpstrait 中scaladsl/Flow.scala因此它同时适用于Source、Flow、SubFlow、SubSource等所有FlowOps的子类型。Java DSL 侧则在 javadsl/Flow.scala 与对应的Source、SubFlow、SubSource中提供内部直接委托给 Scala 实现public alsoTo(that: Graph[SinkShape[Out], _]): javadsl.Flow[In, Out, Mat] new Flow(delegate.alsoTo(that))注意参数类型是Graph[SinkShape[Out], _]而非具体的Sink实现——这意味着你可以传入任何满足SinkShape[Out]的图Sink、SubFlow、Flow的某种连接结果等具备很高的组合灵活性。三、Reactive Streams 语义背压是关键文档用一段 callout 明确给出了alsoTo的 Reactive Streams 语义语义行为emits发射当元素可用且附加Sink与下游同时存在需求demand时backpressures背压当下游或附加Sink背压时completes完成当上游完成时cancels取消当下游或附加Sink取消时这段语义描述在源码注释中有一模一样的表述scaladsl/Flow.scala且与Broadcast的行为完全吻合Broadcast会等待所有输出都具备需求才发射元素因此只要附加Sink处理缓慢整条流都会被背压。这是alsoTo与wireTap最本质的差异见下一节。四、alsoTo vs wireTap背压还是丢弃文档中虽然没有展开对比但源码注释反复强调了一个关键区别It is similar towireTapbut will backpressure instead of dropping elements when the givenSinkis not ready.scaladsl/Flow.scala两者的选择标准非常清晰alsoTo有背压的旁路。附加Sink未就绪时主线流会被迫放慢backpressure保证附加路径不丢失任何元素。适合对数据完整性要求高的场景如事件归档、审计日志、精确计量。wireTap无背压的旁路。附加Sink未就绪时元素被直接丢弃主线流不受影响。适合日志、监控等丢了也无所谓的辅助路径。一句话总结追求零丢失选alsoTo追求主线零干扰选wireTap。代价是alsoTo的吞吐上限受限于最慢的分支。五、变体alsoToAll 与 alsoToMatalsoToAll同时挂接多个 Sink当需要把元素同时发给多个附加Sink时可以使用alsoToAllscaladsl/Flow.scaladef alsoToAll(those: Graph[SinkShape[Out], _]*): Repr[Out]其实现与alsoTo如出一辙只是把Broadcast的输出数扩展为those.size 1端口0留给主下游其余端口分别连接各个Sink。特殊情况下传入空列表时直接返回this原流不产生任何额外开销def alsoToAll(those: Graph[SinkShape[Out], _]*): Repr[Out] those match { case those if those.isEmpty this.asInstanceOf[Repr[Out]] case _ via(GraphDSL.create() { implicit b import GraphDSL.Implicits._ val bcast b.add(BroadcastOut) for ((that, idx) - those.zipWithIndex) bcast.out(idx 1) ~ that FlowShape(bcast.in, bcast.out(0)) }) }测试用例 FlowAlsoToAllSpec.scala 验证了多 Sink 与空参两种形态Source.single(1).alsoToAll(sink1, sink2).runWith(sink3) // 元素同时进入 sink1、sink2、sink3 Source.single(1).alsoToAll().runWith(sink1) // 等价于直接 runWithJava 侧对应alsoToAll(those: Graph[SinkShape[Out], _]*)标注了varargs与SafeVarargs可直接传多个 Sinkjavadsl/Flow.scala。alsoToMat同时获取附加 Sink 的物化值默认情况下alsoTo的物化值就是当前流自身的物化值附加 Sink 的物化值被忽略。若需要同时拿到附加 Sink 的物化结果例如Sink.seq收集到的元素序列使用alsoToMatscaladsl/Flow.scaladef alsoToMatMat2, Mat3(matF: (Mat, Mat2) Mat3): ReprMat[Out, Mat3]测试 FlowFutureFlowSpec.scala 中大量使用了这个形态例如Flow[Int].alsoToMat(Sink.seq)(Keep.right)Keep.right表示最终物化值取附加Sink一侧这里是Future[Seq[Int]]。源码注释建议优先使用内部优化的Keep.left/Keep.right组合器而不是手写透传函数。Java 侧对应alsoToMat(that, matF)接收Function2[Mat, M2, M3]javadsl/Flow.scala。六、实战示例Scala旁路写文件 主线继续处理import akka.actor.ActorSystem import akka.stream.scaladsl.{Flow, Sink, Source} implicit val system: ActorSystem ActorSystem(alsoTo-demo) Source(1 to 100) .alsoTo(Flow[Int].map(i s$i\n).to(Sink.file(...))) // 旁路落盘零丢失 .filter(_ % 2 0) .runWith(Sink.foreach(n println(seven: $n)))Java旁路采集指标并获取物化结果import akka.stream.javadsl.*; FlowInteger, Integer, NotUsed flow Flow.of(Integer.class) .alsoTo(Sink.foreach(n - metrics.record(n))); // 旁路打点背压式保真若附加 Sink 需要快速处理以免拖慢主线可先在旁路上用buffer或async边界隔离但请记住alsoTo的语义决定了任何分支的积压最终都会传导回上游这是与wireTap的本质区别。七、典型应用场景审计与归档主流程处理业务数据的同时把原始元素完整写入事件日志或归档存储alsoTo的背压特性保证审计数据不丢。指标采集与监控旁路发送元素给指标Sink如计数、直方图聚合适合对精度有要求而不仅是采样的场景。数据复制/扇出同一元素同时进入多个下游管道如实时计算 批处理落库alsoToAll可一次挂接多个目标。调试与观测临时挂一个打印Sink观察流经元素无需改动主链路若担心影响吞吐可改用wireTap。八、小结alsoTo用最朴素的方式Broadcast 两个输出实现了流经即旁路的能力是 Akka Streams Fan-out 家族中最易用的成员之一。掌握它的关键在于三点语义上它是背压式旁路区别于丢元素的wireTap、结构上它是Broadcast(2, eagerCancel true)、组合上它有alsoToAll多 Sink与alsoToMat取物化值两个变体。需要零丢失的旁路处理时优先考虑它。参考资源仓库内路径官方文档akka-docs/src/main/paradox/stream/operators/Source-or-Flow/alsoTo.mdFan-out 算子索引akka-docs/src/main/paradox/stream/operators/index.mdScala 实现含alsoTo/alsoToAll/alsoToMatakka-stream/src/main/scala/akka/stream/scaladsl/Flow.scalaJava 实现akka-stream/src/main/scala/akka/stream/javadsl/Flow.scala测试用例FlowAlsoToAllSpec.scala、FlowFutureFlowSpec.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),仅供参考
返回列表