Akka Streams `zipLatestWith` 操作符详解:以最新元素组合多条流
Akka Streams zipLatestWith 操作符详解:以最新元素组合多条流
zipLatestWith 是 Akka Streams 中一个重要的 Fan-in(扇入) 操作符,它将当前流(Source/Flow)与另一条 Source 通过用户提供的 combine 函数组合起来,并且总是取每条输入流的最新元素参与组合。阅读本文后,你将掌握 zipLatestWith 的完整签名、发射/背压/完成/取消语义、eagerComplete 参数的行为差异,以及它在仓库源码与测试中的真实实现细节,从而在实际项目中正确选用它与 zip、zipWith 等相近操作符。
操作符概览
zipLatestWith 的功能可以用一句话概括:当所有输入流都至少产生过一个元素后,只要任意一条输入流出现新元素,就用该新元素与其他输入流各自"最新"的元素一起调用 combine 函数,并把结果推向下游。
该操作符属于文档中归类的 Fan-in operators 一族,其最典型的场景是"合并两个异步数据源,且只关心彼此最新的状态"——例如将传感器实时读数与配置流合并,或把两个高频更新的消息源按最新值对齐。
方法签名
zipLatestWith 同时提供 Scala DSL 与 Java DSL 两个版本,且都可以作用于 Source 与 Flow。
Scala DSL
在 scaladsl/Flow.scala 中定义如下:
def zipLatestWithOut2, Out3(
combine: (Out, Out2) => Out3): Repr[Out3]
// 可显式指定 eagerComplete 的重载版本
def zipLatestWithOut2, Out3(
combine: (Out, Out2) => Out3): Repr[Out3]
参数含义:
that:要与之组合的另一条流,类型为Graph[SourceShape[Out2], _],因此不仅可以直接传Source,也可以传入任意具有SourceShape的图;combine:(Out, Out2) => Out3组合函数,接收当前流的一个元素与that流的一个元素,返回组合结果;eagerComplete(可选,默认true):控制完成语义,详见下文专节。
两个重载版本均通过 via(zipLatestWithGraph(...)) 将当前流接入内部构造的 ZipLatestWith 图阶段,最终返回 Repr[Out3](即原流的 Source/Flow 形态,保留原有类型与物化值)。
此外还提供保留两侧物化值的 zipLatestWithMat(见 scaladsl/Flow.scala):
def zipLatestWithMatOut2, Out3, Mat2, Mat3(
combine: (Out, Out2) => Out3)(
matF: (Mat, Mat2) => Mat3): Repr[Out3]
Java DSL
在 javadsl/Flow.scala 中,Java 版本使用 akka.japi.function.Function2 表达组合函数:
public <Out2, Out3> javadsl.Flow<In, Out3, Mat> zipLatestWith(
Graph<SourceShape<Out2>, ?> that,
function.Function2<Out, Out2, Out3> combine)
public <Out2, Out3> javadsl.Flow<In, Out3, Mat> zipLatestWith(
Graph<SourceShape<Out2>, ?> that,
boolean eagerComplete,
function.Function2<Out, Out2, Out3> combine)
Source 上的用法与之类似,只是返回 javadsl.Source<Out3, Mat>。Source/Flow 的 Java DSL 中同样有 zipLatestWithMat 重载(见 javadsl/Flow.scala)。
工作原理:永远取"最新"而不是"配对"
理解 zipLatestWith 的关键在于它与 zip/zipWith 的本质区别:
zip/zipWith:严格按元素到达次序一一配对(第 N 个元素与第 N 个元素组合),任一输入多出的元素会被"挂起"等待配对;zipLatestWith:不关心元素是否成对到达,始终用每个输入流最近一次到达的元素参与组合。新元素到达哪一侧,就用它替换该侧保存的"最新值",并立即与新值组合后发射。
具体流程(文档与源码 scaladsl/Flow.scala 的注释一致):
- 在所有输入流都至少产生过一个元素之前,不发射任何元素(两侧各需一个"初始最新值");
- 此后,每当任一输入流产生新元素,就取该新元素 + 另一侧保存的最新元素,调用
combine,将结果推向下游; - 若同一时间两侧都有新元素,则合并为一次组合发射。
例如左流依次产生 1, 2,右流依次产生 10, 20, 30,则发射序列为 1+10、2+10、2+20、2+30——注意右侧第二个元素 20 是与左侧"最新值" 2 组合的,这正是"latest"语义的体现。
eagerComplete 参数:决定完成行为
该参数是 zipLatestWith 区别于 zipWith 的又一个重要特性,默认值为 true。从 scaladsl/Flow.scala 可以看到:
def zipLatestWithOut2, Out3(
combine: (Out, Out2) => Out3): Repr[Out3] =
zipLatestWith(that, eagerComplete = true)(combine)
即不带 eagerComplete 参数的版本等价于 eagerComplete = true。两种取值的行为如下(见 scaladsl/Flow.scala 与 scaladsl/Graph.scala 的注释):
| 取值 | 完成语义 |
|---|---|
eagerComplete = true(默认) |
任一上游完成,整个流立即完成,并取消其余上游 |
eagerComplete = false |
等待所有上游都完成后才完成;但若某上游在尚未发出任何元素时就已结束,组合流仍会立即完成 |
图形 DSL 中的完整语义在 scaladsl/Graph.scala 中有更完整的表述:若某些上游在尚未发出任何值之前就已经完成,则组合流立即完成;若所有上游都已产生过值,且 eagerComplete 为 true(默认),则任一上游完成即完成整个流,否则等待全部上游完成。
源码实现剖析
DSL 层到图阶段的组装
Flow.zipLatestWith 并不直接持有运行逻辑,而是通过 zipLatestWithGraph 把两侧输入接到一个 ZipLatestWith 图阶段上(见 scaladsl/Flow.scala):
protected def zipLatestWithGraphOut2, Out3, M(
combine: (Out, Out2) => Out3): Graph[FlowShape[Out @uncheckedVariance, Out3], M] =
GraphDSL.createGraph(that) { implicit b => r =>
val zip = b.add(ZipLatestWithOut, Out2, Out3)
r ~> zip.in1
FlowShape(zip.in0, zip.out)
}
其中 r ~> zip.in1 把参数 that 接入 ZipLatestWith 的 in1 端口,而当前流自身经 via 接入 in0,从而将"一条外部 Source + 当前流"折叠成一个 Flow。
图形 DSL 中的 ZipLatest / ZipLatestWith
在 scaladsl/Graph.scala 中可以同时找到两个相关阶段:
ZipLatest[A, B]:输出(A, B)二元组,其实现是class ZipLatestA, B extends ZipLatestWith2A, B, (A, B),即把Tuple2.apply当作组合函数;默认构造器new ZipLatest()也采用eagerComplete = true;ZipLatestWith[A, B, C]:接受显式的组合函数与eagerComplete参数,二者共享同一个底层ZipLatestWith2阶段实现。
因此从源码结构可以推断:Source.zipLatest(that)(不带 With、输出二元组)本质上是 zipLatestWith 的"元组特化",两者共享同一套最新元素组合逻辑。
测试验证:从用例看语义边界
仓库在 GraphZipLatestWithSpec.scala 中为 ZipLatestWith 提供了完整的语义测试,是理解该操作符行为的最佳佐证:
happy case(L43-L77):左流为 1, 2 后接永不发射的 never 源,右流手动推送 10, 20, 30,得到 11, 12, 22, 32——20 与左侧最新值 2 组合得 22,30 与最新值 2 组合得 32,完整演示了"始终取最新元素"的行为。
sad case(L79-L110):当 combine 内部抛出异常(如除零),流以该异常失败,下游收到 onError,验证异常会沿流传播。
非 eager 版本(L112-L151):eagerComplete = false 时,右侧完成(sendComplete)后流并不结束,左侧继续推送 20、30 仍能与右侧保存的最新值 2 组合得到 22、32,直到左侧也完成才 expectComplete()。
eager 版本(L153-L176):eagerComplete = true 时,左侧 Source(1 to 2) 完成后下游立即 expectComplete(),验证"任一上游完成即完成"。
一方立即/延迟完成或失败的场景(L180-L212):无论哪一侧先完成或先失败,下游都会相应收到完成或错误信号,确认完成/取消传播不受输入侧顺序影响。
此外,DslConsistencySpec.scala 通过反射校验 Scala DSL 与 Java DSL 的方法集合一致性,确保 zipLatestWith 在两种 API 下行为对齐。
Reactive Streams 语义
zipLatestWith 在 Reactive Streams 契约下的行为总结如下(与文档及 scaladsl/Flow.scala 源码注释一致):
| 信号 | 行为 |
|---|---|
| emits(发射) | 所有输入流都至少有一个元素可用之后,每当任一输入流产生新元素就发射一次组合结果 |
| backpressures(背压) | 当下游背压时向上游施加背压 |
| completes(完成) | eagerComplete = true 时任一上游完成即完成;false 时等待所有上游完成 |
| cancels(取消) | 当下游取消时取消整个流 |
实战示例
Scala:合并两个异步源的最新值
import akka.actor.ActorSystem
import akka.stream.scaladsl._
implicit val system: ActorSystem = ActorSystem("zipLatestWith-demo")
val prices: Source[Double, _] = Source(Seq(1.0, 1.5, 2.0))
val rates: Source[Double, _] = Source(Seq(10.0, 11.0))
prices
.zipLatestWith(rates)((price, rate) => price * rate)
.runForeach(println)
Java:使用 Function2 组合
import akka.stream.javadsl.Source;
import akka.japi.function.Function2;
Source.from(java.util.Arrays.asList(1, 2, 3))
.zipLatestWith(
Source.from(java.util.Arrays.asList(10, 20)),
new Function2<Integer, Integer, Integer>() {
public Integer apply(Integer left, Integer right) {
return left + right;
}
});
需要物化值时:zipLatestWithMat
当需要同时拿到两侧流的物化值(例如 actor 引用、Future 完成信号等)时,使用 zipLatestWithMat 并通过 matF 合并两侧 Mat:
val combined: Source[String, (NotUsed, NotUsed)] =
srcA.zipLatestWithMat(srcB)((a, b) => s"$a-$b")(Keep.both)
小结
zipLatestWith 面向"多路流的最新值组合"场景,与严格按序配对的 zipWith 形成互补。选型时可以依据两条准则:
- 需要一一配对、不丢元素时用
zip/zipWith; - 需要始终以各输入的最新值组合、并允许旧值被覆盖时用
zipLatestWith,同时按业务需要选择默认的eagerComplete = true(任一上游完成即完成)或false(等待全部完成)。
其完整实现可从 scaladsl/Flow.scala、scaladsl/Graph.scala 与测试 GraphZipLatestWithSpec.scala 进一步研读。