- 后端
- 并发编程
- 异步编程
【免费下载链接】akka-core
A platform to build and run apps that are elastic, agile, and resilient. SDK, libraries, and hosted environments.
导读
Unzip是 Akka Streams 中一个典型的 Fan-out(扇出)算子:它接收一个由二元组(two element tuples)构成的流,把每个元素的第一个分量与第二个分量分别分发到两个不同的下游流(downstream)。本文以 Unzip 官方算子文档 为主体,结合 akka-stream 模块的源码实现与单元测试,完整讲解它的端口结构、Scala/Java 两种 DSL 的用法、Reactive Streams 背压语义、底层实现原理以及工程实践中的注意事项。读完本文,你将掌握如何用Unzip(及配套的zip)在两个异构类型的流之间做高效、无损的拆分与重组。
一、Unzip 是什么:Fan-out 家族的一员
在 Akka Streams 中,Fan-out 算子拥有一个输入端口和多个输出端口,它们要么把元素路由到不同输出,要么把同一个元素同时发射到多个输出。Unzip属于前者。
在 算子索引 中,Fan-out 家族还包括 Balance(负载均衡扇出)、Broadcast(每个元素广播到 n 个输出)和 Partition(按谓词分派),而Unzip的职责非常单一:
Takes a stream of two element tuples and unzips the two elements into two different downstreams.
即:输入一个由二元组构成的流,把每个二元组的两个元素分别"解开",送往两个不同的下游。它是 Fan-in 算子zip的逆操作,常与zip配对使用,用于把一个携带"数据 + 元数据"或"键 + 值"的复合流拆开,分别进行独立处理后再合并。
二、签名与端口结构
原文档中的 Signature 部分由StreamOperatorsIndexGenerator自动生成(见 project/StreamOperatorsIndexGenerator.scala),其核心签名定义在 DSL 源码中。
Scala DSL
在 scaladsl/Graph.scala 中:
object Unzip { /** Create a new `Unzip`. */ def apply[A, B](): Unzip[A, B] = new Unzip() } final class Unzip[A, B]() extends UnzipWith2(A, B), A, B { override def toString = "Unzip" }- 类型参数
[A, B]分别对应二元组的第一、第二分量类型; Unzip有一个in输入端口和left、right两个输出端口;- 从源码结构看,
Unzip本质上是一个特化的UnzipWith2:它把"拆分函数"固定为恒等函数ConstantFun.scalaIdentityFunction,即直接把(A, B)原样拆成A和B。
Java DSL
在 javadsl/Graph.scala 中:
object Unzip { /** Creates a new `Unzip` operator with the specified output types. */ def create[A, B](): Graph[FanOutShape2[A Pair B, A, B], NotUsed] = UnzipWith.create(ConstantFun.javaIdentityFunction[Pair[A, B]]) def createA, B: Graph[FanOutShape2[A Pair B, A, B], NotUsed] = create[A, B]() }Java 版本返回Graph[FanOutShape2[A Pair B, A, B], NotUsed],输入元素类型是akka.japi.Pair<A, B>;重载的create(left, right)版本只是类型提示,不参与运行时的实际拆分。
三、完整用法示例
Scala:通过 GraphDSL 使用 Unzip
Unzip是纯图形算子(GraphStage),在 算子索引 中明确说明这类算子"目前没有流式(fluent)API 可用,必须借助 Graph DSL 使用"。下面的示例来自仓库测试 GraphUnzipSpec.scala,展示了把Int -> String的元组流拆成两个分支、并分别做不同变换:
import akka.stream.{ ClosedShape, OverflowStrategy } import akka.stream.scaladsl._ RunnableGraph .fromGraph(GraphDSL.create() { implicit b => import GraphDSL.Implicits._ val unzip = b.add(Unzip[Int, String]()) Source(List(1 -> "a", 2 -> "b", 3 -> "c")) ~> unzip.in unzip.out1 ~> Flow[String].buffer(16, OverflowStrategy.backpressure) ~> Sink.ignore unzip.out0 ~> Flow[Int].buffer(16, OverflowStrategy.backpressure).map(_ * 2) ~> Sink.ignore ClosedShape }) .run()关键点:
b.add(Unzip[Int, String]())把算子加入 GraphDSL 构建器;unzip.in、unzip.out0(left)、unzip.out1(right)分别接入上游和两个下游;- 两个输出端口可以接完全不同类型的后续流程(这里是
String分支和Int分支),这正是Unzip相对Broadcast的核心差异——Broadcast的所有输出共享同一元素类型。
Java:通过 GraphDSL 使用 Unzip
Java 版本使用akka.japi.Pair作为输入元素类型,同样需要 GraphDSL:
import akka.japi.Pair; import akka.stream.ClosedShape; import akka.stream.javadsl.*; RunnableGraph.fromGraph( GraphDSL.create(builder -> { FanOutShape2<Pair<Integer, String>, Integer, String> unzip = builder.add(Unzip.create(Integer.class, String.class)); builder.from(Source.from(Arrays.asList( Pair.create(1, "a"), Pair.create(2, "b"), Pair.create(3, "c")))) .to(unzip.in()); builder.from(unzip.out0()).to(Sink.ignore()); builder.from(unzip.out1()).to(Sink.ignore()); return ClosedShape.getInstance(); })) .run(system);四、Reactive Streams 语义(背压行为)
原文档给出了Unzip的官方 Reactive Streams 语义,这也是理解它性能特征的关键:
| 行为 | 触发条件 |
|---|---|
| emits(发射) | 当所有输出端口都停止背压、且上游有可用输入元素时 |
| backpressures(背压) | 当任意一个输出端口背压时 |
| completes(完成) | 当上游完成时 |
这段语义在 scaladsl/Graph.scala 与 javadsl/Graph.scala 的 scaladoc 中完全一致,还额外补充了一条:
Cancels whenany downstream cancels(当任意下游取消时取消)
语义的工程含义
- 发射需要"全部就绪":
Unzip不会为某个更快的下游单独推进,只有当left和right两个下游都愿意接收时,才会消费下一个输入元组。这意味着两个下游的实际吞吐量由较慢的一方决定——它不会为快的一方提前缓冲数据。 - 任一背压即整体背压:如果某个下游处理缓慢(如写入慢速 IO),
Unzip会把背压信号传回上游,从而避免无界缓冲。 - 下游取消的容错:尽管语义上"任意下游取消则取消整个算子",仓库测试 GraphUnzipSpec.scala 验证了
Unzip的FanOut基类实际上会把取消信号隔离——测试 "produce to right downstream even though left downstream cancels" 证明:当 left 下游取消后,right 下游依然能收到全部"a"、"b"、"c"并正常完成。
五、底层实现原理:FanOut 与 TransferPhase
Unzip的运行时实现位于 impl/FanOut.scala,它是 Akka Streams 内部 API(标注@InternalApi private[akka]):
@InternalApi private[akka] class Unzip(attributes: Attributes) extends FanOut(attributes, outputCount = 2) { outputBunch.markAllOutputs() initialPhase( 1, TransferPhase(primaryInputs.NeedsInput && outputBunch.AllOfMarkedOutputs) { () => primaryInputs.dequeueInputElement() match { case (a, b) => outputBunch.enqueue(0, a) outputBunch.enqueue(1, b) case t: akka.japi.Pair[_, _] => outputBunch.enqueue(0, t.first) outputBunch.enqueue(1, t.second) case t => throw new IllegalArgumentException( s"Unable to unzip elements of type ${t.getClass.getName}, " + s"can only handle Tuple2 and akka.japi.Pair!") } }) }从源码结构看,其核心设计可以归纳为三点:
- 继承自
FanOut,固定outputCount = 2:Unzip直接复用 Fan-out 的基础设施(输入子接收器、输出批次管理),无需从零实现背压协调。 outputBunch.markAllOutputs()+AllOfMarkedOutputs:这正是第四节语义的代码级体现——转移阶段(TransferPhase)要求"上游有输入(NeedsInput)且所有被标记的输出都有需求(AllOfMarkedOutputs)"时才消费一个元素,从而保证只有当两个下游都就绪时才发射。- 严格的类型约束:
dequeueInputElement()的返回只接受 ScalaTuple2(case (a, b))和 Javaakka.japi.Pair(case t: akka.japi.Pair[_, _])两种形态;遇到其他类型会抛出IllegalArgumentException,并明确提示 "can only handle Tuple2 and akka.japi.Pair!"。
此外,FanOut基类还实现了故障传播(pumpFailed→fail)、Actor 终止时的清理(postStop中取消输入并向下游发送AbruptTerminationException)以及"不可重启"策略(postRestart直接抛IllegalStateException),保证算子状态机的一致性与背压/取消信号的正确传递。
六、测试验证:行为契约一览
仓库为Unzip提供了完整的契约测试,位于 GraphUnzipSpec.scala,可概括为以下行为保证:
- "unzip to two subscribers":输入
List(1 -> "a", 2 -> "b", 3 -> "c"),left 分支经map(_ * 2)收到2、4、6,right 分支收到"a"、"b"、"c";验证了按元素顺序、按分量类型正确拆分。 - "produce to right downstream even though left downstream cancels"与反向用例:验证单向下游取消不会阻塞另一侧的正常发射与完成。
- 测试基类配置了
akka.stream.materializer.initial-input-buffer-size = 2,并配合TestSubscriber.manualProbe手动控制request(n),精确验证了背压与按需发射的行为。
七、与 UnzipWith、zip 的关系及选型建议
Unzip与 UnzipWith 同属 Fan-out 拆分算子,但适用场景不同:
Unzip:输入必须是二元组(Tuple2 或akka.japi.Pair),拆分方式是固定的恒等拆分,无自定义函数,语义最直观。UnzipWith:输入可以是任意类型,通过用户提供的 splitter 函数把每个元素拆成最多 6 路输出,灵活度更高(Unzip的 DSL 签名extends UnzipWith2(A, B), A, B也印证了二者是同一套机制的特化与泛化关系)。
在流式 DSL 中,Source/Flow上还有与Unzip目标相近的alsoTo、wireTap等旁路算子,但它们属于"主线照常 + 旁路观察"的语义,与Unzip的"一对二独立拆分"并不等价,选型时需要注意区分。
反向操作上,Unzip是 Fan-in 算子 zip 的逆操作:zip把两个流的元素合并为元组,Unzip把元组流拆回两路。典型的组合模式是"zip合 → 联合处理 →Unzip拆"或"Unzip拆 → 并行处理 →zip再合",用于在异构数据流之间做结构化的分离与重组。
八、注意事项总结
- 必须使用 GraphDSL:
Unzip没有流式(fluent)API,只能在GraphDSL.create()中通过b.add(...)使用,参见 stream-graphs.md。 - 输入类型严格:Scala 侧为
(A, B)元组,Java 侧为akka.japi.Pair<A, B>;传入其他类型会触发IllegalArgumentException(见 impl/FanOut.scala)。 - 吞吐由慢下游决定:由于"所有输出就绪才发射",若某个下游长期无需求,整个流会被阻塞。需要为慢分支预留缓冲(如
buffer(16, OverflowStrategy.backpressure))或改用其他策略。 - 两侧类型可不同:
out0(left)与out1(right)分别承载A与B类型,这是与Broadcast的本质区别。 - 完成与取消语义:上游完成则算子完成;任意下游取消时,算子整体取消,但实现层面允许未取消的一侧把已分发元素消费完毕(见测试用例验证)。
通过本文的讲解,你可以放心地在 Akka Streams 图编排中使用Unzip完成"二元组流 → 两路独立流"的拆分,并借助源码级语义理解其背压行为,避免在慢下游场景下踩坑。
- 后端
- 并发编程
- 异步编程
【免费下载链接】akka-core
A platform to build and run apps that are elastic, agile, and resilient. SDK, libraries, and hosted environments.
相关推荐
Akka Streams UnzipWith 算子完全指南:用拆分函数将一个输入流扇出为多个下游
Akka Streams UnzipWith 算子完全指南:用拆分函数将一个输入流扇出为多个下游 本指南围绕 Akka Streams 内置的 Fan out(
后端并发编程异步编程Akka Streams Partition 算子完全指南:按分区函数将流扇出到多个下游
Akka Streams Partition 算子完全指南:按分区函数将流扇出到多个下游 Partition 是 Akka Streams 中一个典型的扇出(F
后端并发编程异步编程Akka Streams Source.zipN 详解:将多个上游源合并为元素序列流
Akka Streams Source.zipN 详解:将多个上游源合并为元素序列流 导读 Source.zipN 是 Akka Streams 中用于多路合并
后端并发编程异步编程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考