Akka Stream:如何保留依赖外部元素的内部Source物化值?
Akka Stream:保留内部源Materialized Value的实现方案
问题描述
我有一个类型为Source<A, Mat1>的源outer,还有一个根据输入A创建内部源的方法Source<B, Mat2> buildInner(A a)。我需要得到一个Source<B, Mat2>类型的结果,但使用flatMapConcat时会丢失Mat2——因为内部源的创建依赖outer的元素,无法提前通过inner.preMaterialize(actorSystem)获取Mat2实例。
解决方案
Akka Stream的flatMap系列算子(flatMapConcat/flatMapMerge)设计上会丢弃内部源的materialized value,仅保留外部源的。要保留Mat2,需要换思路处理:
1. 收集所有内部源的Mat2
如果需要获取所有内部源的Mat2实例,可以结合preMaterialize和Sink来收集:
// 将每个内部源的元素与对应的Mat2绑定 Source<Pair<B, Mat2>, Mat1> elementsWithMat = outer.flatMapConcat(a -> { Pair<Source<B, NotUsed>, Mat2> preMat = buildInner(a).preMaterialize(actorSystem); // 把元素和Mat2打包成Pair输出 return preMat.first().map(b -> Pair.create(b, preMat.second())); }); // 用Sink.fold收集所有Mat2到列表中 Sink<Pair<B, Mat2>, CompletionStage<List<Mat2>>> matCollector = Sink.fold(new ArrayList<>(), (list, pair) -> { list.add(pair.second()); return list; }); // 运行流时同时获取外部源的Mat1和收集到的Mat2列表 Pair<Mat1, CompletionStage<List<Mat2>>> runResult = elementsWithMat.toMat(matCollector, Keep.both()).run(actorSystem); // 后续可通过runResult.second()获取所有Mat2的异步结果
2. 结果流直接返回单个Mat2(仅适用于outer仅发射一个元素)
如果outer只会产生一个元素,那么可以用Promise保存内部源的Mat2,并替换结果流的materialized value:
Promise<Mat2> matPromise = Promise.create(actorSystem); Source<B, Mat2> result = outer .flatMapConcat(a -> { Pair<Source<B, NotUsed>, Mat2> preMat = buildInner(a).preMaterialize(actorSystem); matPromise.complete(preMat.second()); return preMat.first(); }) .mapMaterializedValue(ignoredMat1 -> matPromise.future().toCompletableFuture());
注意:如果outer发射多个元素,Promise会被多次complete,导致异常,仅适用于单元素场景。
3. 手动用Source.queue控制流(复杂场景)
对于更灵活的场景,可以创建一个队列作为结果流,手动订阅outer并处理内部源:
Promise<Mat2> matPromise = Promise.create(actorSystem); // 创建结果队列 SourceQueueWithComplete<B> queue = Source.<B>queue(100, OverflowStrategy.backpressure()) .to(Sink.ignore()) .run(actorSystem); // 订阅outer,处理每个元素对应的内部源 outer.runForeach(a -> { Pair<Source<B, NotUsed>, Mat2> preMat = buildInner(a).preMaterialize(actorSystem); matPromise.complete(preMat.second()); // 多元素场景需改为收集逻辑 // 将内部源的元素推送到队列 preMat.first().runForeach(b -> queue.offer(b), actorSystem); }, actorSystem); // 结果流绑定队列,并将materialized value设为Mat2 Source<B, Mat2> result = Source.fromPublisher(queue) .mapMaterializedValue(ignored -> matPromise.future().toCompletableFuture());
内容的提问来源于stack exchange,提问作者Sigurd
相关产品推荐
相关产品推荐

