【问题标题】:Akka.Net: how to create persistent views through a journal readerAkka.Net:如何通过期刊阅读器创建持久视图
【发布时间】:2018-03-05 12:43:42
【问题描述】:

根据 Akka.Net 文档,不推荐使用 PersistentView,而应使用 PersistenceQuery。 在 ASP.Net Core 2.0 Web-API 应用程序中,我使用 Akka.Net 和事件源。我正在使用带有事件和快照的 SQL Server 插件来实现持久性。 对于持久视图,我想开始使用 PersistenceQuery。当应用程序启动时,会回放事件以恢复参与者的状态。

我已经实现了一个日志阅读器,它接收事件,并使用它来组成一个视图。问题是,我怎样才能知道最后播放的事件已经到达,以便可以保存组合视图(作为一种快照)?我不想在恢复阶段的每个事件之后保存视图。

现在在初始化 ActorSystem 时启动日志阅读器(通过 Startup.cs 调用)。代码如下所示:

private static void InitialiseJournalReader()
{
    // Obtain read journal by plugin id.
    var readJournal = PersistenceQuery.Get(ActorSystem).ReadJournalFor<SqlReadJournal>("akka.persistence.query.myjournal");

    // Materialize stream, consuming events.
    var materializer = ActorMaterializer.Create(ActorSystem);

    var writer = ActorSystem.ActorOf(CreateViewsActor.GetActorProps(), CreateViewsActor.GetActorName());

    // issue query to journal
    Source<EventEnvelope, NotUsed> source = readJournal.CurrentEventsByTag("MyEvents");
    source.RunForeach(envelope => writer.Ask(envelope.Event), materializer);
}

CreateViewsActor 是一个 Actor,它使用消息来创建一个或多个视图。它还必须保存这些视图(目前以 JSON 格式保存到 SQL Server 表中)。

不幸的是,到目前为止,我还没有找到通过期刊阅读器创建持久视图的工作示例。但也许我一直在寻找错误的地方。 到目前为止,我有以下问题:

  1. 是否有任何通过期刊阅读器创建持久视图的工作示例?
  2. CreateViewsActor(或任何负责创建和保存视图的代码)如何知道所有恢复消息都已处理?
  3. 初始化期刊阅读器的最佳位置是什么?

【问题讨论】:

    标签: c# akka-stream akka.net akka.net-persistence


    【解决方案1】:

    阅读日志可用于多种用途。在大多数情况下,用于从事件生成专门的读取视图。然而,这并不一定意味着演员 - 您可以轻松地将视图转换为数据库表中的更新以获得物化视图。

    如果你想将actors与流结合起来,你可以使用Sink.ActorRefSink.ActorRefWithAck,这取决于你是想在你的actor中包含背压还是在全推力模式下工作。示例:

    using (var materializer = system.Materializer())
    {
        var readJournal = PersistenceQuery.Get(system)
            .ReadJournalFor<SqlReadJournal>("akka.persistence.query.my-read-journal");
    
        var writer = asysem.ActorOf(CreateViewsActor.Props(), CreateViewsActor.GetActorName());
    
        readJournal
            .CurrentEventsByTag("MyEvents")
            .Collect(envelope => envelope.Event as MyEvent)
            .RunWith(Sink.ActorRefWithAck<MyEvent>(writer, 
                onInitMessage: CreateViewsActor.Init.Instance,
                ackMessage: CreateViewsActor.Ack.Instance,
                onCompleteMessage: CreateViewsActor.Done.Instance), materializer);
    }
    

    这里 initack 消息是 从参与者发送的 通知流,何时开始发射或只是发射下一个元素给演员(如果有的话)。一旦流完成,最后一个接收器参数(在完成消息上)将发送给参与者

    如有疑问,您可以随时查看官方 Akka.NET 测试(请参阅:12)。

    关于日志阅读器初始化 - 它是一个对象,它的生命周期与参与者系统绑定,因此它可以在与参与者系统相同的位置进行初始化和使用。

    【讨论】:

    • 当我尝试您的代码时,我收到以下错误:[[akka://my-actor-server/user/tenants/tenant-a402c502-fb1d-4e21-8ae9-b0643eaf7ae3/my-view-actor-a402c502-fb1d-4e21-8ae9-b0643eaf7ae3/StreamSupervisor-0/Flow-0-0-unknown-operation#1013508337]] terminated abruptly Cause: Unknown“未知原因”可能是什么?
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2011-07-13
    • 2020-03-08
    • 2015-02-19
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多