【问题标题】:Adding state between operations within akka stream在akka流中的操作之间添加状态
【发布时间】:2021-05-29 11:07:18
【问题描述】:

下面是我用来计算对象列表中数据流平均值的代码:

import akka.NotUsed;
import akka.actor.ActorSystem;
import akka.stream.javadsl.Sink;
import akka.stream.javadsl.Source;

import java.util.ArrayList;
import java.util.Arrays;
import java.util.List;
import java.util.concurrent.CompletionStage;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeoutException;

public class sd001 {

    private static final ActorSystem system = ActorSystem.create("akkassembly");
    private static List<RData> ls = new ArrayList();

    private static class RData {
        private String id;

        public RData(String id){
            this.id = id;
        }

        public List<Integer> getValues(){
            if(this.id.equalsIgnoreCase("1")) {
                return Arrays.asList(1, 2, 3, 4, 5);
            }
            else {
                return Arrays.asList(1, 2, 3);
            }
        }

        public String getId() {
            return this.id;
        }
    }

    final static List<RData> builderFunction() {
        try {
            ls.add(new RData("1"));
            ls.add(new RData("2"));
            ls.add(new RData("3"));
            Thread.sleep(3000);
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
        return ls;
    }

    private static double calculateAverage(List <Integer> marks) {
        return marks.stream()
                .mapToDouble(d -> d)
                .average()
                .orElse(0.0);
    }

    public static void main(String[] args) throws ExecutionException, InterruptedException, TimeoutException {

        final Source<List<RData>, NotUsed> source2 =
                Source.repeat(NotUsed.getInstance()).map(elem -> builderFunction());

                source2.mapConcat(i -> i)
                .groupBy(3 , x -> x.getId())
                .map(v -> calculateAverage(v.getValues()))
                .to(Sink.foreach(x -> System.out.println(x)))
                .run(system);

    }

}

结果输出:

11:55:27.477 [akkassembly-akka.actor.default-dispatcher-4] INFO akka.event.slf4j.Slf4jLogger - Slf4jLogger started
3.0
2.0
2.0
3.0

所以似乎按预期工作。

我使用 groupBy 方法 (https://doc.akka.io/docs/akka/current/stream/stream-substream.html) 按关联的 id 值对项目的 List 进行分组。如何将 id 值添加到输出平均值的阶段,以便不仅输出平均值,还将 id 打印到屏幕上?我指的阶段是:

.to(Sink.foreach(x -> System.out.println(x)))

一种可能的解决方案是修改方法getValues 并创建一个新参数id 并返回除平均值之外的id,这将允许访问println 中的值以获取Sink .这个解决方案似乎过于复杂。看来我需要在mapto 函数之间携带一个额外的状态(在这种情况下为id)?

【问题讨论】:

    标签: java akka akka-stream


    【解决方案1】:

    通常,Akka Streams 中的阶段不共享状态:它们仅在它们之间传递流的元素。因此,在流的各个阶段之间传递状态的唯一通用方法是将状态嵌入到正在传递的元素中。

    在某些情况下,可以使用SourceWithContext/FlowWithContext

    本质上,FlowWithContext 只是一个包含元素和上下文元组的Flow,但优势在于运算符:FlowWithContext 上的大多数运算符将作用于元素而不是元组,允许您专注于您的应用程序逻辑,而不是担心上下文。

    在这种特殊情况下,由于groupBy 正在执行类似于重新排序元素的操作,FlowWithContext 不支持groupBy,因此您必须将 ID 嵌入到流元素中...

    (...除非您想深入到自定义图形阶段的深层,这可能会使将 ID 嵌入流元素的复杂性相形见绌。)

    【讨论】:

    • "you'll have to embed the IDs into the stream elements" ,这与我建议的解决方案的想法是否相同:“一种可能的解决方案是修改方法 getValues 并创建一个新的参数 id 和除了平均值之外还返回 id,这将允许访问 Sink 的 println 中的值。"
    • 是的。我忘记了 Java 是否为此提供了 Pair 类,否则您将不得不自己动手。
    • 您是说使用自定义图形阶段比嵌入 ID 更复杂吗?
    • 是的。根据我的经验,自定义阶段基本上总是最后的手段。
    猜你喜欢
    • 2016-10-20
    • 2016-05-06
    • 1970-01-01
    • 1970-01-01
    • 2011-07-26
    • 2018-09-02
    • 2015-11-28
    • 1970-01-01
    • 2023-01-18
    相关资源
    最近更新 更多