如何基于Orleans构建事件溯源读模型?相关实现疑问
Great question—let’s unpack this thoroughly, since combining Orleans with Event Sourcing (ES) to build read models is a common but tricky scenario that balances decoupling, performance, and consistency.
First: The General Approach to Read Models in Orleans ES
Before diving into push/pull specifics, let’s lay out the core pattern that works for most scenarios:
- Event Contract Standardization: Define a shared interface (e.g.,
IDomainEvent) for all events emitted by journalled grains, including metadata likeEventId,Timestamp,GrainId, andVersion. This makes it easier to handle events uniformly in read models. - Event Routing: Choose a mechanism to get events from journalled grains to read models (push or pull, covered below).
- Idempotent Projection: Read models must handle duplicate events gracefully (use
EventIdor grain-specific event versions as unique keys to skip reprocessing). - Final Consistency: Accept that read models will eventually reflect all events—Orleans and ES are designed for distributed systems where strict consistency isn’t always feasible or necessary.
Push-Based Implementation Schemes
Push is ideal for low-latency read models, where you want updates to reflect almost immediately. Here are two practical approaches:
1. Orleans Streams (Recommended)
Orleans’ built-in stream abstraction is perfect for decoupling journalled grains from read models. You can structure streams in a few ways:
- Per-Grain-Type Streams: Create a stream for each grain type (e.g.,
StreamId.Create("UserEvents", "All")) so read models interested in user events can subscribe to that single stream. - Per-Grain-Instance Streams: For granular control, each grain instance publishes to its own stream (e.g.,
StreamId.Create("OrderEvents", orderId)). Read models can subscribe to specific instances or use wildcards to listen to all instances of a type.
Example Code:
In your journalled grain (after persisting the event):
public async Task ProcessOrder() { var orderCompletedEvent = new OrderCompletedEvent(GrainId.GetPrimaryKey(), DateTime.UtcNow, OrderStatus.Completed); await ApplyEvent(orderCompletedEvent); // Your ES logic to persist the event // Publish to the stream var streamProvider = GetStreamProvider("Default"); var stream = streamProvider.GetStream<IDomainEvent>(StreamId.Create("OrderEvents", "All")); await stream.OnNextAsync(orderCompletedEvent); }
In your read model grain (subscribing on activation):
public override async Task OnActivateAsync() { var streamProvider = GetStreamProvider("Default"); var orderStream = streamProvider.GetStream<IDomainEvent>(StreamId.Create("OrderEvents", "All")); // Subscribe and handle events await orderStream.SubscribeAsync(async (evt, token) => { if (evt is OrderCompletedEvent completedEvt) { // Update your read model (e.g., write to a SQL view or Redis cache) var orderReadModel = await _dbContext.OrderReadModels.FindAsync(completedEvt.OrderId); orderReadModel.Status = completedEvt.Status; orderReadModel.CompletedAt = completedEvt.Timestamp; await _dbContext.SaveChangesAsync(); // Persist the stream token to resume on activation State.LastProcessedToken = token; await WriteStateAsync(); } }); }
2. Direct Grain-to-Read-Model Calls (Tighter Coupling)
If you don’t need full decoupling, you can have journalled grains directly call read model grains after persisting events. This is simpler but creates a dependency between grains.
Example:
// In journalled grain await ApplyEvent(orderCompletedEvent); var orderReadModel = GrainFactory.GetGrain<IOrderReadModelGrain>(orderCompletedEvent.OrderId); await orderReadModel.ApplyOrderCompletedEvent(orderCompletedEvent);
Pull-Based Implementation Schemes
Pull is better for batch-oriented read models, or when you want to control the timing of updates (e.g., nightly reports). Here are two approaches:
1. Polling Journalled Grains
Create a dedicated projection grain that periodically polls all relevant journalled grains for new events. Each journalled grain exposes a method like GetEventsSince(long lastProcessedVersion) to return unprocessed events.
Example:
// Projection grain logic public async Task RunBatchProjection() { var lastProcessedVersions = State.LastProcessedVersions; // Stored in grain state // Get all grain IDs of the target type (use GrainDirectory) var grainDirectory = GrainFactory.GetGrain<IGrainDirectory>(0); var userGrainIds = await grainDirectory.GetGrainIdsOfType(typeof(IUserGrain)); foreach (var grainId in userGrainIds) { var userGrain = GrainFactory.GetGrain<IUserGrain>(grainId); var lastVersion = lastProcessedVersions.GetValueOrDefault(grainId, 0); var newEvents = await userGrain.GetEventsSince(lastVersion); foreach (var evt in newEvents) { await UpdateUserReadModel(evt); // Apply event to read model } // Update the last processed version for this grain if (newEvents.Any()) { lastProcessedVersions[grainId] = newEvents.Max(e => e.Version); } } State.LastProcessedVersions = lastProcessedVersions; await WriteStateAsync(); // Schedule next run RegisterTimer(_ => RunBatchProjection(), null, TimeSpan.FromHours(1), TimeSpan.FromHours(1)); }
2. Directly Querying the Event Store
If your journalled grains persist events to a shared store (e.g., EventStoreDB, SQL Server event table), you can have read models query the store directly instead of going through grains. This is more efficient for large-scale systems, as it avoids hitting individual grains.
Example:
// Read model service logic (could be a grain or external service) public async Task SyncReadModel() { var lastSyncTime = State.LastSyncTime; // Query the event store for all events since last sync var events = await _eventStore.QueryEventsAfter(lastSyncTime); foreach (var evt in events) { await ApplyEventToReadModel(evt); } State.LastSyncTime = DateTime.UtcNow; await WriteStateAsync(); }
Is Cross-Aggregate Read Models an Unreasonable Approach?
Absolutely not—cross-aggregate read models are not only reasonable but often necessary for real-world business scenarios. For example:
- An order history page that combines data from
OrderGrain,UserGrain, andProductGrain. - A dashboard showing metrics across multiple grain types (e.g., total users, active orders, shipped items).
The key is to accept final consistency—events from different aggregates may arrive at the read model out of order, but over time the view will converge to the correct state. If you need stricter consistency, you can use event ordering mechanisms (e.g., stream sequences) or compensate for out-of-order events in your projection logic.
Final Key Takeaways
- Prefer Orleans Streams for push-based read models—they’re built into Orleans, handle scaling, and decouple producers from consumers.
- Use direct event store queries for pull-based models when possible—it’s more efficient than polling grains.
- Always handle idempotency—duplicate events are inevitable in distributed systems.
- Embrace final consistency—it’s the tradeoff for scalability and decoupling in Orleans + ES systems.
内容的提问来源于stack exchange,提问作者Ilya I. Margolin

