【问题标题】:Akka stateless actors with Persistent Mailboxes具有持久邮箱的 Akka 无状态演员
【发布时间】:2019-11-24 04:42:50
【问题描述】:

我想创建一个包含 1000 多个演员的 Akka 集群。每个 Actor 接收一条消息,进行一些计算并将结果写入一个专用的 Kafka 主题。

它应该部署在集群中,例如 Kubernetes。

我的理解是——如果actor被终止(JVM崩溃、重新部署或其他任何原因),那么其邮箱的内容——以及当前正在处理的消息——都会丢失!

这在我的情况下是完全不可接受的,因此我想实现一种拥有持久邮箱的方法。请注意,参与者本身是无状态的,他们不需要重播消息或重建状态。我所需要的只是在演员被终止时不会丢失消息。

问题是:推荐的方法是什么? Herehere 他们建议实施持久性演员。但就像我说的,我不需要坚持和恢复演员的任何状态。我应该实现基于持久存储(如 SQL 数据库)的自定义邮箱吗?

我还看到在某些版本之前 Akka 支持“持久”邮箱,这似乎是我所需要的。但由于某种原因,他们删除了它,这令人困惑......

【问题讨论】:

    标签: akka actor akka-cluster akka-persistence


    【解决方案1】:

    您可以使用 Kafka 来实现您想要的。 Kafka 主题是持久的(如果您将 Kafka 中的保留设置为 forever 或在主题上启用日志压缩,那么数据将“一直”保留,或者您可以在 Kafka 之外存储偏移量)。

    使用 Akka Streams,您将在广播您生成的消息(关于生成主题)之后提交您收到的消息(关于接收主题),从而为您提供“至少一次”传递语义。(对于“exactly-once”,您可以查看 Kafka Transactions)

    这是来自 Alpakka Kafka 文档的示例:

    Consumer.DrainingControl<Done> control =
        Consumer.committableSource(consumerSettings, Subscriptions.topics(topic))
            .map(
                msg ->
                    ProducerMessage.single(
                        new ProducerRecord<>(targetTopic, msg.record().key(), msg.record().value()),
                        msg.committableOffset() // the passThrough
                        ))
            .via(Producer.flexiFlow(producerSettings))
            .map(m -> m.passThrough())
            .toMat(Committer.sink(committerSettings), Keep.both())
            .mapMaterializedValue(Consumer::createDrainingControl)
            .run(materializer);
    

    您可以通过几种方式将其与(集群池)Actors 集成。最简单的方法是使用Ask 模式。在这种情况下,流会将消息传递给必须在预定义时间内回复的参与者(可能是self())。当收到回复时,它将在提交原始消息之前在目标流上广播。

    这看起来像:

    Consumer.DrainingControl<Done> control =
        Consumer.committableSource(consumerSettings, Subscriptions.topics(topic))
                .mapAsync(1, msg -> 
                    Patterns.ask(actor, msg, Duration.ofSeconds(5))
                        .thenApply(done ->
                            ProducerMessage.single(
                                    new ProducerRecord<>(targetTopic, done.key(), done.value()),
                                    msg.committableOffset() // the passThrough
                            )
                        )
                )
                .via(Producer.flexiFlow(producerSettings))
                .map(m -> m.passThrough())
                .toMat(Committer.sink(committerSettings), Keep.both())
                .mapMaterializedValue(Consumer::createDrainingControl)
                .run(materializer);
    

    如果您有多个可以同时处理消息的参与者,您还可以增加 mapAsync 调用的并行度。

    【讨论】:

      【解决方案2】:

      客户端上使用持久性参与者是针对此类需求的建议。我知道您说您的接收参与者不需要持久性/有状态,但是通过在客户端上使用持久性,您可以在接收参与者终止时重试,或者使用开箱即用的保证消息传递功能来确保它得到处理.本质上,持久性(在客户端)用于持久化所发出的请求,以便 客户端 可以在必要时重新发送消息以“重建邮箱”。

      使用客户端持久化是:

      • 比持久邮箱更高效
      • 防止出现更多故障情况(例如在网络层丢失消息、应用程序逻辑故障)
      • 更灵活,支持更多类型的恢复(例如:只需要恢复部分消息的场景)

      这就是为什么从 Akka 中删除持久邮箱的原因:Akka Persistence/Guaranteed At Least Once 在所有方面都比持久邮箱更好。

      stikkos 对使用 Kafka 的回答也是可行的。我只是担心引入 Kafka 会增加很多复杂性。当然,任何持久性存储都会增加复杂性,所以我想这取决于您已经拥有的内容。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2019-04-23
        • 2017-11-30
        • 2012-04-24
        • 1970-01-01
        • 2012-12-17
        • 2023-03-03
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多