You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何合并两个Akka.Net Persistence事件溯源系统的事件

问题:合并Akka.Net Persistence双数据库事件并按时间戳排序

需要合并两个由Akka.Net Persistence处理的事件溯源系统中的事件,要求按事件时间戳排序。已确认Akka.Stream的MergeSorted算子符合需求(已用数字列表测试,且编写了自定义EventEnvelopComparer)。

当前方案创建了两个ActorSystem:readsystem1连接PostGres数据库db1,readsystem2连接db2。但遇到问题:使用MergeSorted算子时,若基于readsystem1创建ActorMaterializer,仅能加载db1的事件;若基于readsystem2创建则仅加载db2的事件,无法同时加载两者。

测试代码如下:

var actorMaterializer1 = ActorMaterializer.Create(
    readSystem1,
    ActorMaterializerSettings.Create(readSystem1).WithDebugLogging(true));
var readJournal1 = PersistenceQuery.Get(readSystem1)
    .ReadJournalFor<SqlReadJournal>(SqlReadJournal.Identifier);
var source1 = readJournal1.CurrentEventsByPersistenceId("mypersistenceId", 0L, long.MaxValue);
await source1
    .Select(x => ByteString.FromString($"{x.Timestamp.ToString()}{Environment.NewLine}"))
    .RunWith(FileIO.ToFile(new FileInfo(@"c:\tmp\timestamps1.txt")), actorMaterializer1);

// 仅创建materializer就改变了source加载的事件!!!
var actorMaterializer2 = ActorMaterializer.Create(
    readSystem2, 
    ActorMaterializerSettings.Create(readSystem1).WithDebugLogging(true));
var readJournal2 = PersistenceQuery.Get(readSystem2)
    .ReadJournalFor<SqlReadJournal>(SqlReadJournal.Identifier);
var source2 = readJournal2.CurrentEventsByPersistenceId("mypersistenceId", 0L, long.MaxValue);
await source2
    .Select(x => ByteString.FromString($"{x.Timestamp.ToString()}{Environment.NewLine}"))
    .RunWith(FileIO.ToFile(new FileInfo(@"c:\tmp\timestamps2.txt")), actorMaterializer2);

// 使用actorMaterializer1运行,仅能加载db1的事件并自我合并
var source = source1.MergeSorted(source2, new EventEnvelopComparer());
await source
    .Select(x => ByteString.FromString($"{x.Timestamp.ToString()}{Environment.NewLine}"))
    .RunWith(FileIO.ToFile(new FileInfo(@"c:\tmp\timestamps.txt")), actorMaterializer1);

请问如何实现同时加载并合并双库事件?是否可在同一ActorSystem中读取不同数据库或不同事件溯源表?ActorMaterializer是否有解决方案?当前思路是否完全错误?


解决方案

核心问题原因

你的当前思路问题在于:每个Source(来自SqlReadJournal)绑定到创建它的ActorSystem的内部资源(调度器、数据库连接池等)。当你用其中一个ActorSystem的ActorMaterializer运行合并流时,另一个ActorSystem的Source无法获取自身所需的资源,因此不会产生任何事件。


方案一:同一ActorSystem配置多个SqlReadJournal(推荐)

Akka.Net允许在同一个ActorSystem中注册多个独立配置的SqlReadJournal实例,只需给每个实例分配唯一的标识符,分别指向不同的数据库。这是最简洁的解决方案。

步骤1:配置多个ReadJournal

在HOCON配置文件中添加两个SqlReadJournal的配置,分别指定不同的数据库连接、表名(如果需要)和标识符:

akka.persistence.query.journal.sql-db1 {
  class = "Akka.Persistence.Sql.Query.SqlReadJournal, Akka.Persistence.Sql"
  plugin-dispatcher = "akka.actor.default-dispatcher"
  connection-string = "Server=localhost;Port=5432;Database=db1;User Id=postgres;Password=postgres;"
  provider-name = "Npgsql"
  table-name = "event_journal"
  schema-name = "public"
}

akka.persistence.query.journal.sql-db2 {
  class = "Akka.Persistence.Sql.Query.SqlReadJournal, Akka.Persistence.Sql"
  plugin-dispatcher = "akka.actor.default-dispatcher"
  connection-string = "Server=localhost;Port=5432;Database=db2;User Id=postgres;Password=postgres;"
  provider-name = "Npgsql"
  table-name = "event_journal"
  schema-name = "public"
}

步骤2:在代码中获取两个ReadJournal

创建一个ActorSystem,然后分别通过自定义的标识符获取两个SqlReadJournal实例:

var system = ActorSystem.Create("MergeEventSystem");
var materializer = ActorMaterializer.Create(system, ActorMaterializerSettings.Create(system).WithDebugLogging(true));

// 获取db1的ReadJournal
var readJournal1 = PersistenceQuery.Get(system)
    .ReadJournalFor<SqlReadJournal>("akka.persistence.query.journal.sql-db1");
var source1 = readJournal1.CurrentEventsByPersistenceId("mypersistenceId", 0L, long.MaxValue);

// 获取db2的ReadJournal
var readJournal2 = PersistenceQuery.Get(system)
    .ReadJournalFor<SqlReadJournal>("akka.persistence.query.journal.sql-db2");
var source2 = readJournal2.CurrentEventsByPersistenceId("mypersistenceId", 0L, long.MaxValue);

// 合并并输出到文件
var mergedSource = source1.MergeSorted(source2, new EventEnvelopComparer());
await mergedSource
    .Select(x => ByteString.FromString($"{x.Timestamp.ToString()}{Environment.NewLine}"))
    .RunWith(FileIO.ToFile(new FileInfo(@"c:\tmp\timestamps.txt")), materializer);

这种方式下,两个Source都属于同一个ActorSystem,使用同一个Materializer就能正常工作,无需跨系统资源调用。


方案二:跨ActorSystem的StreamRef桥接(备选)

如果必须保留两个独立的ActorSystem,可以使用Akka.Stream的StreamRef功能,将其中一个ActorSystem的Source转换成可跨系统传输的流引用,然后在另一个ActorSystem中消费该引用,实现合并。

步骤1:在源ActorSystem中暴露StreamRef

// 在readSystem2中创建并暴露Source的StreamRef
var source2 = readJournal2.CurrentEventsByPersistenceId("mypersistenceId", 0L, long.MaxValue);
var streamRef = source2.RunWith(StreamRefs.SourceRef<EventEnvelope>().ToSink(), actorMaterializer2);
var sourceRef = await streamRef;

步骤2:在目标ActorSystem中消费StreamRef并合并

// 在readSystem1中获取sourceRef并转换成本地Source
var remoteSource = sourceRef.Source<EventEnvelope>(system);

// 合并source1和remoteSource
var mergedSource = source1.MergeSorted(remoteSource, new EventEnvelopComparer());
await mergedSource
    .Select(x => ByteString.FromString($"{x.Timestamp.ToString()}{Environment.NewLine}"))
    .RunWith(FileIO.ToFile(new FileInfo(@"c:\tmp\timestamps.txt")), actorMaterializer1);

此方案需要确保两个ActorSystem之间能正确通信(配置Akka.Remote或Akka.Cluster),相对复杂,仅在无法使用单一ActorSystem的场景下考虑。


总结

你的初始思路(双ActorSystem直接合并)的核心问题是跨系统资源隔离导致的Source无法正常触发。推荐优先使用方案一,通过同一ActorSystem配置多ReadJournal的方式实现,代码更简洁、性能更优。

内容的提问来源于stack exchange,提问作者Daniele Armanasco

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.18 12:50:26