EventSourcing订阅通用实践及GetEventStore中聚合事件重放的疑问
Hey there, let's break down your questions step by step based on my hands-on experience with EventStore and Event Sourcing!
Should you create separate subscriptions for each observer?
Absolutely—this is a much cleaner, scalable approach compared to a single global subscription. Here's why:
- Targeted replay: When you need to rebuild a specific Read Model, you only need to restart its dedicated subscription (or directly read its relevant streams) instead of filtering through every event in the entire store.
- Isolation: If one Read Model has an issue (like a bug in event handling), you can take it offline for fixes without disrupting other parts of your system.
- Performance: Global subscriptions force you to filter every event for every Read Model, which gets slow as your event volume grows. Separate subscriptions let each observer only process events it cares about.
The only caveat is managing multiple subscriptions—EventStore handles this well as long as you use a shared connection (or connection pool) instead of spawning a new connection per subscription.
How to get the Stream ID?
Stream IDs are your responsibility to define when you write events to EventStore. A common pattern is to use a consistent naming convention like:
{AggregateType}-{AggregateId}
For example: Order-123e4567-e89b-12d3-a456-426614174000 or Customer-JohnDoe-123.
When handling events (either in a global subscription or a stream-specific one), you can access the Stream ID directly from the resolved event:
void OnEventMaterialized(EventStoreConnection connection, ResolvedEvent resolvedEvent) { string streamId = resolvedEvent.Event.EventStreamId; // Use streamId to determine if this event belongs to the observer's target aggregate }
If you need to replay events for a specific aggregate, just construct the Stream ID using your naming convention—no need to "discover" it from a list.
Do you need to persist all created Stream IDs?
Nope! You don't have to maintain your own list of Stream IDs. EventStore provides built-in ways to work with streams:
- Read directly by constructed Stream ID: As mentioned above, if you know the aggregate type and ID, you can build the Stream ID and read its events directly with
ReadStreamEventsForwardAsync. - Subscribe to the
$streamssystem stream: This system stream emits an event every time a new stream is created. You can subscribe to it to track new streams if needed, but for most replay scenarios, constructing the Stream ID is sufficient.
Code Examples
1. Subscribe to a single aggregate stream for a specific Read Model
public async Task SubscribeToAggregateStream<TReadModel>( TReadModel readModel, string aggregateType, Guid aggregateId, EventStoreConnection connection) { var streamId = $"{aggregateType}-{aggregateId}"; // Load the last saved position for this Read Model (if you persisted it) var lastPosition = await _readModelPositionRepository.GetLastPositionAsync( readModel.GetType().Name, streamId); var subscription = await connection.SubscribeToStreamAsync( streamId, lastPosition ?? StreamPosition.Start, (sub, resolvedEvent) => { // Pass the event to your Read Model handler readModel.Handle(resolvedEvent.Event); // Persist the new position after handling _readModelPositionRepository.SavePositionAsync( readModel.GetType().Name, streamId, resolvedEvent.Event.EventNumber); }, (sub, reason, ex) => { Console.WriteLine($"Subscription to {streamId} dropped: {reason}. Error: {ex?.Message}"); // Add retry logic here if needed } ); }
2. Replay all events for a single aggregate to rebuild a Read Model
If you don't want to use a subscription (for one-off rebuilds), you can directly read the stream:
public async Task ReplayAggregateEvents<TReadModel>( TReadModel readModel, string aggregateType, Guid aggregateId, EventStoreConnection connection) { var streamId = $"{aggregateType}-{aggregateId}"; var currentPosition = StreamPosition.Start; const int batchSize = 100; // Adjust based on your event size while (true) { var result = await connection.ReadStreamEventsForwardAsync( streamId, currentPosition, batchSize, resolveLinkTos: true); foreach (var resolvedEvent in result.Events) { readModel.Handle(resolvedEvent.Event); } if (result.IsEndOfStream) break; currentPosition = result.NextEventNumber; } // Save the final position so future subscriptions don't reprocess everything await _readModelPositionRepository.SavePositionAsync( readModel.GetType().Name, streamId, currentPosition); }
Key Event Sourcing Takeaways
- Decouple Read Models: Each Read Model should operate independently—this makes maintenance, scaling, and debugging far easier.
- Persist subscription positions: Always track where each Read Model left off, so you don't have to reprocess every event on every restart.
- Leverage EventStore's built-in tools: System streams like
$alland$streamseliminate the need to manage your own stream inventories.
内容的提问来源于stack exchange,提问作者James Woodley

