如何合并两个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

