You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.14 23:31:04