【问题标题】:How to maintain Alpakka/Akka Streams source state across application restarts?如何在应用程序重新启动时维护 Alpakka/Akka Streams 源状态?
【发布时间】:2019-12-27 17:16:52
【问题描述】:

我是 Alpakka 的新手,正在考虑将其用于系统集成。在应用程序重新启动时维护 Akka Streams 源状态的理想方法是什么?

例如:假设我正在使用以下内容来连续读取一些输入数据并将其转储到某个地方。如果它运行了 4 小时,然后整个 JVM 崩溃并重新启动(例如 k8s 重新启动我的 pod 左右)怎么办:

someSource
    .via(someTransformation)
    .via(someOtherTransformation)
    .toMap(...)
    .run()

我了解,如果 someSource 是 Kafka 源或 Kinesis 源或其他一些有状态的源,他们可以跟踪其偏移量或检查点,并或多或少地从中断的地方重新启动。

但是,许多其他来源没有这样的概念,例如Cassandra 源、文件源或 RDBMs 源。例如,如果我关闭并重新启动rdms example中提供的代码,它每次都会从顶部重新启动。

我是否正确理解没有开箱即用的机制来解决这个问题,s.t.我们必须手动处理吗?我会想象这个功能会很受欢迎,以至于它会以某种方式处理。如果不是,人们通常如何解决这个问题?您是否使用 Akka 持久性将一些游标存储在几个演员中?或者您是否将原点偏移与输出数据一起存储并在启动时重新读取?

还是我看错了?

【问题讨论】:

    标签: scala akka akka-stream alpakka


    【解决方案1】:

    由于您建议的原因,这是一个非常普遍需要的功能。

    然而,实现这一点的唯一通用、可靠的方法是使用 akka 持久性,这可能是 Akka 生态系统中最重的(例如,它需要选择数据库)依赖项。除此之外,它将在某种程度上特定于源。有些(例如 Kafka、Kinesis)有一种方法可以满足几乎所有场景的需求,但对于其他人来说,如何存储消费状态的细节会有很多差异的意见。一般而言,Akka 和 Alpakka 倾向于回避意见。

    【讨论】:

    • 感谢李维斯的洞察力。好的,我想这是有道理的,跨技术实现统一的方式会很棘手。我认为我将采用的方法是将源状态存储为目标数据的一部分。例如,我正在从一些 HTTP 流端点读取数据并将记录推送到 kafka。我可以简单地在存储源“光标”的 kafka 记录中设置一个标题,重新启动后我可以轻松地重新发现最新成功处理的光标。
    猜你喜欢
    • 2019-06-12
    • 1970-01-01
    • 1970-01-01
    • 2010-12-13
    • 1970-01-01
    • 2018-01-14
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多