使用 MediatR 实现事件驱动架构
📋 概述
事件驱动架构(Event-Driven Architecture, EDA)是一种松耦合的分布式系统设计模式。MediatR 的通知机制(INotification)为在 .NET 应用中实现事件驱动提供了简洁而强大的基础。
学习目标
- ✅ 理解事件驱动架构的核心概念
- ✅ 掌握 MediatR 通知机制的使用
- ✅ 学习领域事件的实现模式
- ✅ 实践事件发布/订阅模式
- ✅ 了解事件溯源和最终一致性
🎯 核心概念
什么是事件驱动架构?
事件驱动架构是一种软件架构模式,其中系统的组件通过产生和消费事件来进行通信。
mermaid
graph LR
A[事件生产者] -->|发布事件| B[事件总线]
B -->|分发给| C[事件消费者1]
B -->|分发给| C2[事件消费者2]
B -->|分发给| C3[事件消费者3]
style A fill:#e1f5ff
style B fill:#fff4e1
style C fill:#e8f5e9
style C2 fill:#e8f5e9
style C3 fill:#e8f5e9核心要素
| 要素 | 说明 | 示例 |
|---|---|---|
| 事件(Event) | 表示已发生的事实 | OrderCreatedEvent |
| 事件生产者 | 触发事件的组件 | 订单服务 |
| 事件消费者 | 处理事件的组件 | 通知服务、库存服务 |
| 事件总线 | 事件分发机制 | MediatR IMediator |
🏗️ 架构设计
整体架构图
mermaid
graph TB
subgraph "API 层"
API[OrdersController]
end
subgraph "应用层 - 命令"
CreateCmd[CreateOrderCommand]
end
subgraph "应用层 - 处理器"
Handler[CreateOrderHandler]
end
subgraph "领域层"
Order[Order Aggregate]
Events[Domain Events]
end
subgraph "基础设施层"
DbContext[DbContext]
end
subgraph "事件处理器"
EmailHandler[EmailNotificationHandler]
InventoryHandler[InventoryUpdateHandler]
AnalyticsHandler[AnalyticsHandler]
CacheHandler[CacheInvalidationHandler]
end
API --> CreateCmd
CreateCmd --> Handler
Handler --> Order
Handler --> DbContext
Order -->|触发| Events
Events -->|发布到| MediatR
MediatR --> EmailHandler
MediatR --> InventoryHandler
MediatR --> AnalyticsHandler
MediatR --> CacheHandler💻 实战示例:电商订单系统
场景描述
当用户下单后,需要执行以下操作:
- 创建订单(核心业务)
- 发送确认邮件(通知服务)
- 扣减库存(库存服务)
- 记录分析数据(分析服务)
- 清除缓存(缓存服务)
步骤 1:定义事件
csharp
// Domain/Events/OrderCreatedEvent.cs
using MediatR;
namespace OrderManagement.Domain.Events;
/// <summary>
/// 订单创建事件
/// </summary>
public record OrderCreatedEvent : INotification
{
public int OrderId { get; init; }
public Guid OrderNumber { get; init; }
public int CustomerId { get; init; }
public decimal TotalAmount { get; init; }
public List<OrderItemSnapshot> Items { get; init; } = new();
public DateTime CreatedAt { get; init; }
}
public record OrderItemSnapshot
{
public int ProductId { get; init; }
public string ProductName { get; init; }
public int Quantity { get; init; }
public decimal UnitPrice { get; init; }
}csharp
// Domain/Events/OrderCancelledEvent.cs
public record OrderCancelledEvent : INotification
{
public int OrderId { get; init; }
public Guid OrderNumber { get; init; }
public string Reason { get; init; }
public DateTime CancelledAt { get; init; }
}csharp
// Domain/Events/OrderShippedEvent.cs
public record OrderShippedEvent : INotification
{
public int OrderId { get; init; }
public Guid TrackingNumber { get; init; }
public string Carrier { get; init; }
public DateTime ShippedAt { get; init; }
}步骤 2:在聚合根中触发事件
csharp
// Domain/AggregatesModel/OrderAggregate/Order.cs
using OrderManagement.Domain.Events;
namespace OrderManagement.Domain.AggregatesModel.OrderAggregate;
public class Order : IAggregateRoot
{
private readonly List<INotification> _domainEvents = new();
public IReadOnlyCollection<INotification> DomainEvents => _domainEvents.AsReadOnly();
public int Id { get; private set; }
public Guid OrderNumber { get; private set; }
public int CustomerId { get; private set; }
public OrderStatus Status { get; private set; }
// ... 其他属性 ...
public Order(int customerId, Address shippingAddress, IEnumerable<OrderItem> items)
{
CustomerId = customerId;
ShippingAddress = shippingAddress;
OrderNumber = Guid.NewGuid();
Status = OrderStatus.Pending;
CreatedAt = DateTime.UtcNow;
_orderItems.AddRange(items);
RecalculateTotalAmount();
// 触发领域事件
AddDomainEvent(new OrderCreatedEvent
{
OrderId = Id,
OrderNumber = OrderNumber,
CustomerId = CustomerId,
TotalAmount = TotalAmount.Amount,
Items = _orderItems.Select(item => new OrderItemSnapshot
{
ProductId = item.ProductId,
ProductName = item.ProductName,
Quantity = item.Quantity,
UnitPrice = item.UnitPrice.Amount
}).ToList(),
CreatedAt = CreatedAt
});
}
public void Cancel(string reason)
{
if (Status == OrderStatus.Cancelled)
throw new DomainException("订单已取消");
if (Status == OrderStatus.Delivered)
throw new DomainException("订单已送达,无法取消");
Status = OrderStatus.Cancelled;
UpdatedAt = DateTime.UtcNow;
// 触发订单取消事件
AddDomainEvent(new OrderCancelledEvent
{
OrderId = Id,
OrderNumber = OrderNumber,
Reason = reason,
CancelledAt = DateTime.UtcNow
});
}
protected void AddDomainEvent(INotification eventItem)
{
_domainEvents.Add(eventItem);
}
public void ClearDomainEvents()
{
_domainEvents.Clear();
}
}步骤 3:实现事件处理器
3.1 邮件通知处理器
csharp
// Application/EventHandlers/OrderCreatedEmailHandler.cs
using MediatR;
using Microsoft.Extensions.Logging;
using OrderManagement.Domain.Events;
using OrderManagement.Application.Interfaces;
namespace OrderManagement.Application.EventHandlers;
/// <summary>
/// 订单创建邮件通知处理器
/// </summary>
public class OrderCreatedEmailHandler : INotificationHandler<OrderCreatedEvent>
{
private readonly IEmailService _emailService;
private readonly ILogger<OrderCreatedEmailHandler> _logger;
public OrderCreatedEmailHandler(
IEmailService emailService,
ILogger<OrderCreatedEmailHandler> logger)
{
_emailService = emailService;
_logger = logger;
}
public async Task Handle(OrderCreatedEvent notification, CancellationToken cancellationToken)
{
_logger.LogInformation("开始发送订单确认邮件,订单ID: {OrderId}", notification.OrderId);
try
{
var emailMessage = new EmailMessage
{
To = await GetCustomerEmail(notification.CustomerId),
Subject = $"订单确认 - {notification.OrderNumber}",
Body = GenerateOrderConfirmationEmail(notification)
};
await _emailService.SendAsync(emailMessage, cancellationToken);
_logger.LogInformation("订单确认邮件发送成功,订单ID: {OrderId}", notification.OrderId);
}
catch (Exception ex)
{
_logger.LogError(ex, "订单确认邮件发送失败,订单ID: {OrderId}", notification.OrderId);
// 注意:这里不抛出异常,避免影响主流程
}
}
private async Task<string> GetCustomerEmail(int customerId)
{
// 从数据库或缓存获取客户邮箱
return "customer@example.com";
}
private string GenerateOrderConfirmationEmail(OrderCreatedEvent orderEvent)
{
return $@"
<h2>订单确认</h2>
<p>尊敬的客户,您的订单已创建成功!</p>
<ul>
<li>订单号: {orderEvent.OrderNumber}</li>
<li>订单金额: ¥{orderEvent.TotalAmount:F2}</li>
<li>商品数量: {orderEvent.Items.Count}</li>
</ul>
<p>我们将尽快为您发货。</p>
";
}
}3.2 库存更新处理器
csharp
// Application/EventHandlers/InventoryUpdateHandler.cs
using MediatR;
using Microsoft.Extensions.Logging;
using OrderManagement.Domain.Events;
using OrderManagement.Application.Interfaces;
namespace OrderManagement.Application.EventHandlers;
/// <summary>
/// 库存更新处理器
/// </summary>
public class InventoryUpdateHandler : INotificationHandler<OrderCreatedEvent>
{
private readonly IInventoryService _inventoryService;
private readonly ILogger<InventoryUpdateHandler> _logger;
public InventoryUpdateHandler(
IInventoryService inventoryService,
ILogger<InventoryUpdateHandler> logger)
{
_inventoryService = inventoryService;
_logger = logger;
}
public async Task Handle(OrderCreatedEvent notification, CancellationToken cancellationToken)
{
_logger.LogInformation("开始扣减库存,订单ID: {OrderId}", notification.OrderId);
foreach (var item in notification.Items)
{
try
{
await _inventoryService.ReserveStockAsync(
item.ProductId,
item.Quantity,
cancellationToken);
_logger.LogInformation(
"商品 {ProductId} 库存预留成功,数量: {Quantity}",
item.ProductId,
item.Quantity);
}
catch (InsufficientStockException ex)
{
_logger.LogError(
ex,
"商品 {ProductId} 库存不足,请求数量: {Quantity}",
item.ProductId,
item.Quantity);
// 可以触发补偿操作(如取消订单)
await TriggerCompensatingTransaction(notification.OrderId, item);
}
}
}
private async Task TriggerCompensatingTransaction(int orderId, OrderItemSnapshot item)
{
// 触发补偿事务,如取消订单
_logger.LogWarning("触发补偿事务,订单ID: {OrderId}", orderId);
// 实现逻辑...
}
}3.3 数据分析处理器
csharp
// Application/EventHandlers/AnalyticsHandler.cs
using MediatR;
using Microsoft.Extensions.Logging;
using OrderManagement.Domain.Events;
using OrderManagement.Application.Interfaces;
namespace OrderManagement.Application.EventHandlers;
/// <summary>
/// 数据分析处理器
/// </summary>
public class AnalyticsHandler : INotificationHandler<OrderCreatedEvent>
{
private readonly IAnalyticsService _analyticsService;
private readonly ILogger<AnalyticsHandler> _logger;
public AnalyticsHandler(
IAnalyticsService analyticsService,
ILogger<AnalyticsHandler> logger)
{
_analyticsService = analyticsService;
_logger = logger;
}
public async Task Handle(OrderCreatedEvent notification, CancellationToken cancellationToken)
{
_logger.LogInformation("记录订单分析数据,订单ID: {OrderId}", notification.OrderId);
try
{
var analyticsEvent = new AnalyticsEvent
{
EventType = "OrderCreated",
Timestamp = notification.CreatedAt,
Properties = new Dictionary<string, object>
{
["OrderId"] = notification.OrderId,
["CustomerId"] = notification.CustomerId,
["TotalAmount"] = notification.TotalAmount,
["ItemCount"] = notification.Items.Count,
["Products"] = notification.Items.Select(i => i.ProductId).ToList()
}
};
await _analyticsService.TrackEventAsync(analyticsEvent, cancellationToken);
_logger.LogInformation("订单分析数据记录成功,订单ID: {OrderId}", notification.OrderId);
}
catch (Exception ex)
{
_logger.LogError(ex, "订单分析数据记录失败,订单ID: {OrderId}", notification.OrderId);
}
}
}3.4 缓存失效处理器
csharp
// Application/EventHandlers/CacheInvalidationHandler.cs
using MediatR;
using Microsoft.Extensions.Caching.Distributed;
using Microsoft.Extensions.Logging;
using OrderManagement.Domain.Events;
using System.Text.Json;
namespace OrderManagement.Application.EventHandlers;
/// <summary>
/// 缓存失效处理器
/// </summary>
public class CacheInvalidationHandler : INotificationHandler<OrderCreatedEvent>
{
private readonly IDistributedCache _cache;
private readonly ILogger<CacheInvalidationHandler> _logger;
public CacheInvalidationHandler(
IDistributedCache cache,
ILogger<CacheInvalidationHandler> logger)
{
_cache = cache;
_logger = logger;
}
public async Task Handle(OrderCreatedEvent notification, CancellationToken cancellationToken)
{
_logger.LogInformation("清除相关缓存,订单ID: {OrderId}", notification.OrderId);
try
{
// 清除客户订单列表缓存
var customerOrdersCacheKey = $"customer:{notification.CustomerId}:orders";
await _cache.RemoveAsync(customerOrdersCacheKey, cancellationToken);
// 清除产品销量缓存
foreach (var item in notification.Items)
{
var productSalesCacheKey = $"product:{item.ProductId}:sales";
await _cache.RemoveAsync(productSalesCacheKey, cancellationToken);
}
_logger.LogInformation("缓存清除成功,订单ID: {OrderId}", notification.OrderId);
}
catch (Exception ex)
{
_logger.LogError(ex, "缓存清除失败,订单ID: {OrderId}", notification.OrderId);
}
}
}步骤 4:发布事件
csharp
// Application/Commands/Orders/CreateOrder/CreateOrderHandler.cs
using MediatR;
using OrderManagement.Application.Interfaces;
using OrderManagement.Domain.AggregatesModel.OrderAggregate;
namespace OrderManagement.Application.Commands.Orders.CreateOrder;
public class CreateOrderHandler : IRequestHandler<CreateOrderCommand, int>
{
private readonly IOrderRepository _orderRepository;
private readonly ICustomerRepository _customerRepository;
private readonly IUnitOfWork _unitOfWork;
private readonly IMediator _mediator;
public CreateOrderHandler(
IOrderRepository orderRepository,
ICustomerRepository customerRepository,
IUnitOfWork unitOfWork,
IMediator mediator)
{
_orderRepository = orderRepository;
_customerRepository = customerRepository;
_unitOfWork = unitOfWork;
_mediator = mediator;
}
public async Task<int> Handle(CreateOrderCommand request, CancellationToken cancellationToken)
{
// 1. 验证客户
var customer = await _customerRepository.GetByIdAsync(request.CustomerId, cancellationToken);
if (customer == null)
throw new ApplicationException($"客户 {request.CustomerId} 不存在");
// 2. 创建订单项
var orderItems = request.Items.Select(item =>
new OrderItem(item.ProductId, item.ProductName, item.Quantity, Money.FromDecimal(item.UnitPrice))
).ToList();
// 3. 创建订单(会触发领域事件)
var order = new Order(request.CustomerId, request.ShippingAddress.ToDomain(), orderItems);
// 4. 持久化
await _orderRepository.AddAsync(order, cancellationToken);
await _unitOfWork.SaveChangesAsync(cancellationToken);
// 5. 发布领域事件
var domainEvents = order.DomainEvents.ToList();
order.ClearDomainEvents();
foreach (var domainEvent in domainEvents)
{
await _mediator.Publish(domainEvent, cancellationToken);
}
// 6. 返回订单ID
return order.Id;
}
}🔧 高级特性
1. 事件过滤
只处理符合条件的事件:
csharp
// Application/EventHandlers/VipOrderCreatedHandler.cs
public class VipOrderCreatedHandler : INotificationHandler<OrderCreatedEvent>
{
private readonly IVipService _vipService;
public async Task Handle(OrderCreatedEvent notification, CancellationToken cancellationToken)
{
// 只处理 VIP 客户的订单
if (!await _vipService.IsVipCustomerAsync(notification.CustomerId))
return;
// VIP 订单特殊处理
await ProcessVipOrder(notification);
}
}2. 事件重试机制
csharp
// Application/Behaviors/EventRetryBehavior.cs
using MediatR;
using Polly;
namespace OrderManagement.Application.Behaviors;
public class EventRetryBehavior<TNotification> : IPipelineBehavior<TNotification, Unit>
where TNotification : INotification
{
private readonly ILogger<EventRetryBehavior<TNotification>> _logger;
public EventRetryBehavior(ILogger<EventRetryBehavior<TNotification>> logger)
{
_logger = logger;
}
public async Task<Unit> Handle(
TNotification notification,
RequestHandlerDelegate<Unit> next,
CancellationToken cancellationToken)
{
var policy = Policy
.Handle<Exception>()
.WaitAndRetryAsync(
retryCount: 3,
sleepDurationProvider: attempt => TimeSpan.FromSeconds(Math.Pow(2, attempt)),
onRetry: (exception, timeSpan, retryCount, context) =>
{
_logger.LogWarning(exception,
"事件处理失败,第 {RetryCount} 次重试: {EventType}",
retryCount,
typeof(TNotification).Name);
});
await policy.ExecuteAsync(() => next());
return Unit.Value;
}
}3. 事件顺序保证
csharp
// 按优先级执行事件处理器
public class PriorityOrderEventHandler : INotificationHandler<OrderCreatedEvent>
{
public int Order => 1; // 优先级最高
public async Task Handle(OrderCreatedEvent notification, CancellationToken cancellationToken)
{
// 优先处理
}
}
// 注册时指定顺序
services.AddTransient<INotificationHandler<OrderCreatedEvent>, PriorityOrderEventHandler>();
services.AddTransient<INotificationHandler<OrderCreatedEvent>, NormalOrderEventHandler>();4. 异步事件总线
csharp
// Infrastructure/EventBus/InMemoryEventBus.cs
using MediatR;
using Microsoft.Extensions.DependencyInjection;
namespace OrderManagement.Infrastructure.EventBus;
public class InMemoryEventBus : IEventBus
{
private readonly IServiceProvider _serviceProvider;
public InMemoryEventBus(IServiceProvider serviceProvider)
{
_serviceProvider = serviceProvider;
}
public async Task PublishAsync<TEvent>(TEvent @event, CancellationToken cancellationToken = default)
where TEvent : INotification
{
using var scope = _serviceProvider.CreateScope();
var mediator = scope.ServiceProvider.GetRequiredService<IMediator>();
await mediator.Publish(@event, cancellationToken);
}
}🎨 最佳实践
1. 事件命名规范
✅ 推荐做法:
- 使用过去时态:
OrderCreated、OrderUpdated、OrderCancelled - 明确表示已发生的事实
- 使用领域语言
❌ 避免做法:
- 使用现在时态:
CreateOrder、UpdateOrder - 使用命令式命名:
CreateOrderCommand
2. 事件粒度
✅ 推荐做法:
- 细粒度事件:
OrderCreated、OrderShipped、OrderDelivered - 每个事件代表一个明确的业务事实
- 便于订阅者选择性订阅
❌ 避免做法:
- 粗粒度事件:
OrderChanged(包含所有变更) - 事件过大(包含过多数据)
3. 错误处理
✅ 推荐做法:
csharp
public async Task Handle(OrderCreatedEvent notification, CancellationToken cancellationToken)
{
try
{
// 处理逻辑
}
catch (Exception ex)
{
_logger.LogError(ex, "事件处理失败: {EventType}", typeof(TEvent).Name);
// 不抛出异常,避免影响主流程
}
}❌ 避免做法:
csharp
public async Task Handle(OrderCreatedEvent notification, CancellationToken cancellationToken)
{
// 直接抛出异常会影响主流程和其他事件处理器
await DoSomething(); // 可能抛出异常
}4. 幂等性保证
csharp
public class IdempotentEventHandler : INotificationHandler<OrderCreatedEvent>
{
private readonly IDistributedCache _cache;
public async Task Handle(OrderCreatedEvent notification, CancellationToken cancellationToken)
{
var processedKey = $"processed_event:{notification.OrderId}:{notification.GetType().Name}";
// 检查是否已处理
if (await _cache.GetStringAsync(processedKey, cancellationToken) != null)
{
return; // 已处理,跳过
}
// 处理事件
await ProcessEvent(notification);
// 标记为已处理(过期时间 24 小时)
await _cache.SetStringAsync(
processedKey,
"true",
new DistributedCacheEntryOptions { AbsoluteExpirationRelativeToNow = TimeSpan.FromHours(24) },
cancellationToken);
}
}5. 性能优化
csharp
// 批量处理事件
public class BatchEventHandler : INotificationHandler<OrderCreatedEvent>
{
private readonly ConcurrentQueue<OrderCreatedEvent> _eventQueue = new();
private readonly Timer _flushTimer;
public BatchEventHandler()
{
// 每 5 秒批量处理一次
_flushTimer = new Timer(_ => FlushEvents(), null, TimeSpan.Zero, TimeSpan.FromSeconds(5));
}
public Task Handle(OrderCreatedEvent notification, CancellationToken cancellationToken)
{
_eventQueue.Enqueue(notification);
return Task.CompletedTask;
}
private void FlushEvents()
{
var events = new List<OrderCreatedEvent>();
while (_eventQueue.TryDequeue(out var evt))
{
events.Add(evt);
}
if (events.Any())
{
// 批量处理
ProcessBatch(events);
}
}
}🔍 调试技巧
1. 事件追踪
csharp
// Application/Diagnostics/EventTracingBehavior.cs
public class EventTracingBehavior<TNotification> : IPipelineBehavior<TNotification, Unit>
where TNotification : INotification
{
private readonly ILogger<EventTracingBehavior<TNotification>> _logger;
private readonly Stopwatch _stopwatch;
public async Task<Unit> Handle(
TNotification notification,
RequestHandlerDelegate<Unit> next,
CancellationToken cancellationToken)
{
_stopwatch.Start();
_logger.LogInformation("开始处理事件: {EventType}", typeof(TNotification).Name);
try
{
return await next();
}
finally
{
_stopwatch.Stop();
_logger.LogInformation(
"事件处理完成: {EventType}, 耗时: {ElapsedMs}ms",
typeof(TNotification).Name,
_stopwatch.ElapsedMilliseconds);
}
}
}2. 事件日志
csharp
// Infrastructure/Diagnostics/EventLogger.cs
public class EventLogger
{
private readonly ILogger<EventLogger> _logger;
public async Task LogEventAsync<TEvent>(TEvent @event, string status, Exception? exception = null)
where TEvent : INotification
{
var logEntry = new EventLogEntry
{
EventType = typeof(TEvent).Name,
Timestamp = DateTime.UtcNow,
Status = status,
Payload = JsonSerializer.Serialize(@event),
Exception = exception?.ToString()
};
await SaveToDatabaseAsync(logEntry);
if (exception != null)
{
_logger.LogError(exception, "事件处理失败: {@LogEntry}", logEntry);
}
else
{
_logger.LogInformation("事件处理成功: {@LogEntry}", logEntry);
}
}
}📊 监控与可观测性
1. 指标收集
csharp
// Infrastructure/Metrics/EventMetrics.cs
using System.Diagnostics.Metrics;
public class EventMetrics
{
private static readonly Meter Meter = new("OrderManagement.Events");
private static readonly Counter<long> EventsProcessed = Meter.CreateCounter<long>("events.processed");
private static readonly Histogram<double> EventDuration = Meter.CreateHistogram<double>("events.duration.milliseconds");
public static void RecordEventProcessed(string eventType)
{
EventsProcessed.Add(1, new KeyValuePair<string, object?>("event.type", eventType));
}
public static void RecordEventDuration(string eventType, double durationMs)
{
EventDuration.Record(durationMs, new KeyValuePair<string, object?>("event.type", eventType));
}
}2. OpenTelemetry 集成
csharp
// Infrastructure/Observability/EventTracing.cs
using OpenTelemetry.Trace;
public class EventTracing
{
private static readonly ActivitySource ActivitySource = new("OrderManagement.Events");
public static async Task TraceEventAsync<TEvent>(Func<Task> handler) where TEvent : INotification
{
using var activity = ActivitySource.StartActivity($"Process {typeof(TEvent).Name}");
try
{
await handler();
activity?.SetStatus(Status.Ok);
}
catch (Exception ex)
{
activity?.SetStatus(Status.Error.WithDescription(ex.Message));
throw;
}
}
}🚀 扩展:事件溯源
什么是事件溯源?
事件溯源(Event Sourcing)是一种将状态变化存储为事件序列的模式。
mermaid
graph LR
A[命令] --> B[聚合根]
B --> C[生成事件]
C --> D[事件存储]
D --> E[投影到读模型]
E --> F[查询]简单实现
csharp
// Domain/EventSourcing/EventStore.cs
public class EventStore : IEventStore
{
private readonly Dictionary<Guid, List<StoredEvent>> _events = new();
public async Task SaveEventsAsync(Guid aggregateId, IEnumerable<IDomainEvent> events)
{
if (!_events.ContainsKey(aggregateId))
_events[aggregateId] = new List<StoredEvent>();
var storedEvents = events.Select(e => new StoredEvent
{
AggregateId = aggregateId,
EventType = e.GetType().Name,
Data = JsonSerializer.Serialize(e),
Timestamp = DateTime.UtcNow
});
_events[aggregateId].AddRange(storedEvents);
await PersistToDatabaseAsync();
}
public async Task<IEnumerable<IDomainEvent>> GetEventsAsync(Guid aggregateId)
{
if (!_events.ContainsKey(aggregateId))
return Enumerable.Empty<IDomainEvent>();
return _events[aggregateId].Select(ReconstructEvent);
}
private IDomainEvent ReconstructEvent(StoredEvent storedEvent)
{
var eventType = Type.GetType(storedEvent.EventType);
return (IDomainEvent)JsonSerializer.Deserialize(storedEvent.Data, eventType);
}
}🎯 总结
通过 MediatR 实现事件驱动架构的优势:
✅ 优势
- 松耦合:生产者和消费者互不知道对方
- 可扩展:轻松添加新的事件处理器
- 灵活性:支持同步和异步处理
- 可靠性:支持重试和补偿事务
- 可测试:每个处理器可独立测试
⚠️ 注意事项
- 最终一致性:事件处理是异步的,可能存在延迟
- 幂等性:确保事件可以被安全地重复处理
- 错误处理:妥善处理失败,避免级联故障
- 监控:建立完善的日志、指标和追踪
📈 适用场景
- ✅ 复杂业务流程
- ✅ 多系统集成
- ✅ 审计需求
- ✅ 实时数据处理
- ❌ 简单的 CRUD 应用
- ❌ 强一致性要求