【问题标题】:Changing source data for akka streams更改 akka 流的源数据
【发布时间】:2021-04-10 23:23:12
【问题描述】:

我正在学习 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 来处理不断变化的源数据。

【问题讨论】:

    标签: java scala akka akka-stream


    【解决方案1】:

    Akka Streams 可以并且经常会一直运行到(不久之前)您的应用停止。例如,通常有一个流消费(例如,使用来自 Alpakka Kafka 的 Kafka 消费者源)Kafka 记录在应用程序中很早就开始,并且在应用程序被终止之前不会停止。

    详细来说,流一直运行到以下时间:

    • 阶段信号完成(例如,在您的示例中,Source.from 将在发出 10 后发出完成信号)
    • 阶段失败(通常引发异常)

    一个对动态数据有用的示例源(不引入 Alpakka 或 Akka HTTP)是 Source.queue,它具体化为一个队列,其中排队的元素可用于流。

    【讨论】:

      猜你喜欢
      • 2019-06-25
      • 1970-01-01
      • 2021-07-06
      • 1970-01-01
      • 1970-01-01
      • 2016-08-22
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多