【问题标题】:Akka Persistence Query and actor shardingAkka Persistence Query 和 actor 分片
【发布时间】:2016-07-04 18:02:12
【问题描述】:

我正在做一个 CQRS Akka 演员应用程序的查询端。

查询参与者设置为集群分片,并填充来自一个持久性查询流的事件。

我的问题是:

  1. 如果集群分片中的参与者之一重新启动,如何恢复?

    • 关闭整个集群分片并回复所有事件?
    • 使集群分片中的参与者成为持久参与者并仅为查询端保存一组新事件?
  2. 如果使用 Persistence Query 填充的 actor 重启,如何取消当前 PQ 并重新启动它?

【问题讨论】:

  • 您是否将查询参与者的状态仅保存在内存中?对于我的查询方面,我使用持久性查询来更新数据库视图。
  • 是的,我只是将状态保存在actor的内存中。
  • 您的视图是否只消耗一个或多个 persitence id?
  • 例如 BookingSaved 事件将被持久化,然后 ApartmentAvailableActor 将使用它,BookedPeriodsActor 也将使用它。所以现在每个查询参与者只会消耗一个persistenceId :)
  • 对我来说,这看起来不像一个 persistentId。假设 Bookings 在您的命令端是一个聚合,每个预订都应该有自己的 persistenceId。所以如果 ApartmentAvailableActor 消费了所有的 BookingSaved 事件,那么它会消费不同的 persistenceIds

标签: akka cqrs akka-stream akka-cluster akka-persistence


【解决方案1】:

如前所述,我将评估将您的查询端持久保存在数据库中。

如果这不是一个选项,并且您想坚持每个分片的单个持久性查询,请在您的查询参与者中执行以下操作:

var inRecovery: Boolean = true;

override def preStart( ) = {
    //Subscribe to your event live stream now, so you don't miss anything during recovery
    // e.g. send Subscription message to your persistence query actor

    //Re-Read everything up to now for recovery
    readJournal.currentEventsByPersistenceId("persistenceId")
        .watchTermination()((_, f) => f pipeTo self) // Send Done to self after recovery is finished
        .map(Replay.apply) // Mark your replay messages
        .runWith( Sink.actorRef( self, tag ) ) // Send all replay events to self
}

override def receive = {
    case Done => // Recovery is finished
        inRecovery = false
        unstashAll() // unstash all normal messages received during recovery

    case Replay( payload ) =>
        //handle replayed messages

    case events: Event =>
        //handle normal events from your persistence query
        inRecovery match {
            case true => stash() // stash normal messages until recovery is done
            case false => 
                // recovery is done, start handling normal events
        }
}


case class Replay( payload: AnyRef )

所以基本上在actor开始之前订阅持久性查询actor并使用所有过去事件的有限流恢复状态,所有事件都通过后终止。在恢复期间存储所有传入事件,这些事件不是重播事件。然后在恢复完成后,取消存储所有内容并开始处理正常消息。

【讨论】:

    猜你喜欢
    • 2016-11-09
    • 1970-01-01
    • 1970-01-01
    • 2021-01-30
    • 2023-04-02
    • 1970-01-01
    • 2012-06-23
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多