【发布时间】: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 .这个解决方案似乎过于复杂。看来我需要在map 和to 函数之间携带一个额外的状态(在这种情况下为id)?
【问题讨论】:
标签: java akka akka-stream