java - 更改 akka 流的源数据
问题描述
我正在学习 Java Akka 流并使用https://doc.akka.io/docs/akka/current/stream/stream-flows-and-basics.html定义了以下内容:
import java.util.Arrays;
import java.util.List;
import java.util.concurrent.CompletionStage;
import java.util.concurrent.ExecutionException;
public class SourceExample {
static ActorSystem system = ActorSystem.create("SourceExample");
public static void main(String args[]) throws ExecutionException, InterruptedException {
final List<Integer> sourceData = Arrays.asList(1, 2, 3, 4, 5, 6, 7, 8, 9, 10);
final Source<Integer, NotUsed> source =
Source.from(sourceData);
final Sink<Integer, CompletionStage<Integer>> sink =
Sink.<Integer, Integer>fold(0, (agg, next) -> agg + next);
final CompletionStage<Integer> sum = source.runWith(sink, system);
System.out.println(sum.toCompletableFuture().get());
}
}
运行此代码按预期运行。
Akka Streams 正在解决的问题是该代码可以重复执行吗?
在现实世界的场景中,sourceData
不会是静态的,Akka Streams 是否对如何处理变化的数据有意见,还是由开发人员决定?
在最简单的情况下,只需在源数据更改时每 X 分钟重新执行一次流式流(例如使用计划任务)。还是 Akka 流长期存在,源数据发生变化,流计算根据某些参数重新执行?
Akka Streams 文档定义了多个数据源,但我不明白应该如何利用 Akka Streams 来处理不断变化的源数据。
解决方案
Akka Streams 可以并且经常运行,直到(不久之前)您的应用程序停止。例如,通常有一个流消费(例如,使用来自 Alpakka Kafka 的 Kafka 消费者源)Kafka 记录在应用程序中很早就开始,并且在应用程序被终止之前不会停止。
详细地说,一个流一直运行到这样的时间:
- 阶段信号完成(例如,在您的示例中,
Source.from
将在发射后发出完成信号10
) - 阶段失败(通常抛出异常)
一个对动态数据有用的示例源(不引入 Alpakka 或 Akka HTTP)是Source.queue
作为队列实现的,对于该队列,入队的元素可用于流。
推荐阅读
- python - openmdao 可以跨 Matlab ExternalCodeComp 计算偏导数而不显式定义它们吗?
- html - 在 Vue.js + TypeScript 中更改图像源属性
- wordpress - 最近从 calendars.icloud.com 读取失败...为什么?
- python - 如何使用参数调用存储为变量的函数 - Python 3
- postgresql - AWS 链接数据库并在多个数据库上运行查询
- r - 从 R 以编程方式启动 BAPI?
- javascript - 我可以在 if/else 条件语句中引用之前定义的函数吗?
- kotlin - 如何在任何条件下使用 assertThat?
- sql - 更新 SQL 表中数据选择中的列并保存
- google-chrome - 是否可以从我的网络应用程序配置 Chrome 的“拨打电话”弹出窗口?