Akka Streams中Keep.left/Keep.right的使用为何未改变流处理结果?
嘿,我完全懂你的困惑——你把物化值(materialized value)和流的元素处理逻辑搞混啦!这是Akka Streams新手最容易踩的坑之一,咱们把这个问题拆解透:
核心概念澄清:物化值≠元素处理逻辑
首先要明确一个关键事实:viaMat的Keep.left/Keep.right绝对不会改变流中元素的处理路径。它的唯一作用是:当你把Source和Flow(或Flow和Flow)组合时,选择保留哪个阶段的物化值作为组合后流的物化值,和元素会不会经过Flow的处理完全无关!
举个直白的类比:
- 流的元素处理就像快递的运输路线:你选了走北京→上海→广州,快递就一定会按这个路线走。
- 物化值就像每个运输节点给你的回执单:
Keep.left是保留北京节点的回执,Keep.right是保留上海节点的回执,但快递还是会走完整个路线。
你的代码为什么结果全一致?
咱们逐个看你的代码片段:
1. 第一个基础示例
val someFuture: NotUsed = numSource.via(incrementFlow).via(doubleFlow).to(Sink.foreach(println)).run
元素路径:1-10 → +1 → *10 → 打印,输出20,30,...110,这完全符合预期。
2. 第二个viaMat示例
val someOtherFuture: Future[Done] = numSource.viaMat(incrementFlow)(Keep.left).viaMat(doubleFlow)(Keep.right).toMat(Sink.foreach(println))(Keep.right).run()
这里的Keep.left/Keep.right只是在选择物化值:
numSource.viaMat(incrementFlow)(Keep.left):组合后的流的物化值是numSource的物化值(也就是NotUsed,因为Source(1 to 10)没有有用的物化值),但元素还是会经过incrementFlow的+1处理。- 后续的
viaMat(doubleFlow)(Keep.right):选择保留doubleFlow的物化值(依然是NotUsed),元素还是会经过*10处理。 - 最后
toMat(Sink.foreach(println))(Keep.right):选择保留Sink的物化值(Future[Done],表示流完成的信号)。
所以元素的处理路径和第一个示例完全一样,自然输出相同的结果。
3. 第三个全Keep.right示例
val someOtherFuture2: Future[Done] = numSource.viaMat(incrementFlow)(Keep.right).viaMat(doubleFlow)(Keep.right).toMat(Sink.foreach(println))(Keep.right).run()
和第二个示例逻辑一致:所有Keep.right只是选择保留每个后续阶段的物化值,但元素依然会依次经过+1和*10的处理,结果当然和前两个一样。
什么时候Keep.left/right才会体现差异?
只有当你的Source/Flow/Sink有有实际意义的物化值时,Keep.left/Keep.right的选择才会影响你拿到的对象。比如:
- 用
Source.queue的物化值是一个队列控制对象,可以用来向流中手动添加元素; - 用
Sink.head的物化值是一个Future[T],可以拿到流中第一个元素; - 用
Flow.groupedWeighted的物化值可能是一个统计计数器。
举个实际的例子:
// 一个Sink,它的物化值是Future[Int](拿到流的最后一个元素) val lastElementSink = Sink.last[Int] // 组合流时选择保留Source的物化值(NotUsed)还是Sink的物化值(Future[Int]) val sourceMat: NotUsed = numSource.toMat(lastElementSink)(Keep.left).run() val sinkMat: Future[Int] = numSource.toMat(lastElementSink)(Keep.right).run() // sinkMat可以用来获取流的最后一个元素 sinkMat.foreach(last => println(s"Last element is: $last")) // 会打印10
这里Keep.left/Keep.right的差异就体现出来了,但流的元素处理(把所有元素传到Sink)并没有改变。
如果你想跳过某个Flow的处理怎么办?
如果你的需求是条件性地跳过某个Flow的处理,那你需要用流的分支逻辑,比如:
- 用
filter过滤掉不需要处理的元素; - 用
flatMapConcat根据元素选择不同的Flow; - 用
GraphDSL构建更复杂的流拓扑。
比如,如果你想让奇数走incrementFlow,偶数直接走doubleFlow,可以这么写:
val branchedFlow = Flow[Int].flatMapConcat { num => if (num % 2 == 1) Source.single(num).via(incrementFlow) else Source.single(num) }.via(doubleFlow) numSource.via(branchedFlow).to(Sink.foreach(println)).run()
这样奇数会先+1再10,偶数直接10,输出就会不一样了。
内容的提问来源于stack exchange,提问作者nikabu

