EventStore gRPC读取$all流:重启后如何正确续读后续事件?
解决EventStore $all流续读未知后续位置事件的问题
核心原因
EventStore的全局位置(Position)是唯一且非连续的,每个事件对应一个精确的Position值,不存在的中间位置会导致ReadAllAsync抛出异常,因此不能直接用已知事件的Position+1作为起始位置。
解决方案
方案1:从已知事件的Position开始读取并跳过该事件
直接以事件B的Position作为起始位置调用ReadAllAsync,读取结果中会包含事件B本身,只需过滤掉该事件即可获取后续的事件A及其他新事件:
// eventBPosition 是已知的事件B的全局位置 var startPosition = new Position(eventBPosition, eventBPosition); // 读取$all流,从事件B的位置开始 var readStream = _client.ReadAllAsync(Direction.Forwards, startPosition, long.MaxValue); await foreach (var resolvedEvent in readStream) { // 跳过事件B本身 if (resolvedEvent.Position.CommitPosition == eventBPosition) { continue; } // 处理后续事件(包括事件A) ProcessEvent(resolvedEvent); }
方案2:使用持久订阅(推荐)
如果是长期需要续读$all流的场景,EventStore的持久订阅是更可靠的方案。持久订阅会自动跟踪消费进度,重启后直接从上次确认的位置继续读取,无需手动管理Position:
- 先创建持久订阅(只需执行一次):
var subscriptionSettings = new PersistentSubscriptionSettings { // 根据需求配置,比如是否自动确认、重试策略等 StartFrom = Position.Start, MaxRetryCount = 3, AckTimeout = TimeSpan.FromSeconds(30) }; // 创建针对$all流的持久订阅,订阅组名自定义 await _client.CreatePersistentSubscriptionAsync( SystemStreams.AllStream, "my-app-all-subscription", subscriptionSettings);
- 订阅并处理事件:
// 订阅持久订阅,自动从上次消费位置续读 var subscription = await _client.SubscribeToPersistentSubscriptionAsync( SystemStreams.AllStream, "my-app-all-subscription", async (sub, resolvedEvent, retryCount, ct) => { // 处理事件 ProcessEvent(resolvedEvent); // 手动确认事件已处理,持久订阅会记录该位置 await sub.Ack(resolvedEvent); }, (sub, reason, ex) => { // 处理订阅断开或异常情况 Console.WriteLine($"订阅断开: {reason}, 异常: {ex?.Message}"); });
使用持久订阅的优势:无需手动存储和管理Position,EventStore会自动维护消费进度,即使应用重启也能无缝续读,同时提供重试、死信队列等可靠性特性。
内容的提问来源于stack exchange,提问作者Maxim Kitsenko
相关产品推荐
相关产品推荐

