Akka Streams `zipLatestWith` 操作符详解:以最新元素组合多条流

原创2026-09-22 19:42:491,357 阅读
文章标签:后端并发编程异步编程

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 的注释一致):

  1. 在所有输入流都至少产生过一个元素之前,不发射任何元素(两侧各需一个"初始最新值");
  2. 此后,每当任一输入流产生新元素,就取该新元素 + 另一侧保存的最新元素,调用 combine,将结果推向下游;
  3. 若同一时间两侧都有新元素,则合并为一次组合发射。

例如左流依次产生 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 形成互补。选型时可以依据两条准则:

  1. 需要一一配对、不丢元素时用 zip/zipWith;
  2. 需要始终以各输入的最新值组合、并允许旧值被覆盖时用 zipLatestWith,同时按业务需要选择默认的 eagerComplete = true(任一上游完成即完成)或 false(等待全部完成)。

其完整实现可从 scaladsl/Flow.scala、scaladsl/Graph.scala 与测试 GraphZipLatestWithSpec.scala 进一步研读。

登录后查看全文
akka-core