Skip to content

使用 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. 创建订单(核心业务)
  2. 发送确认邮件(通知服务)
  3. 扣减库存(库存服务)
  4. 记录分析数据(分析服务)
  5. 清除缓存(缓存服务)

步骤 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 实现事件驱动架构的优势:

✅ 优势 ​

  1. 松耦合:生产者和消费者互不知道对方
  2. 可扩展:轻松添加新的事件处理器
  3. 灵活性:支持同步和异步处理
  4. 可靠性:支持重试和补偿事务
  5. 可测试:每个处理器可独立测试

⚠️ 注意事项 ​

  1. 最终一致性:事件处理是异步的,可能存在延迟
  2. 幂等性:确保事件可以被安全地重复处理
  3. 错误处理:妥善处理失败,避免级联故障
  4. 监控:建立完善的日志、指标和追踪

📈 适用场景 ​

  • ✅ 复杂业务流程
  • ✅ 多系统集成
  • ✅ 审计需求
  • ✅ 实时数据处理
  • ❌ 简单的 CRUD 应用
  • ❌ 强一致性要求

📚 参考资源 ​

Released under the CC BY-SA 4.0 License.