Skip to content

通知(Notification)—— 一对多广播 ​

📢 实现发布/订阅模式,解耦领域事件


📖 概述 ​

MediatR 的通知机制实现了发布/订阅(Publish/Subscribe)模式,允许一个事件触发多个处理器的执行,是实现领域事件和最终一致性的核心工具。

核心特点 ​

  • ✅ 一对多通信:一个通知可以有多个处理器
  • ✅ 完全解耦:发布者不知道订阅者
  • ✅ 开闭原则:新增处理器无需修改发布者
  • ✅ 异步执行:支持并行处理

🎯 INotification 与 INotificationHandler<T> ​

定义通知 ​

csharp
// 通知是一个简单的 POCO 类
public class OrderCreatedNotification : INotification
{
    public Guid OrderId { get; set; }
    public string CustomerEmail { get; set; }
    public decimal TotalAmount { get; set; }
    public DateTime CreatedAt { get; set; }
}

定义多个处理器 ​

csharp
// 处理器 1:发送确认邮件
public class SendConfirmationEmailHandler : INotificationHandler<OrderCreatedNotification>
{
    private readonly IEmailService _emailService;
    private readonly ILogger<SendConfirmationEmailHandler> _logger;

    public SendConfirmationEmailHandler(IEmailService emailService, ILogger<SendConfirmationEmailHandler> logger)
    {
        _emailService = emailService;
        _logger = logger;
    }

    public async Task Handle(OrderCreatedNotification notification, CancellationToken cancellationToken)
    {
        _logger.LogInformation("发送订单确认邮件: {OrderId}", notification.OrderId);
        
        await _emailService.SendAsync(
            notification.CustomerEmail,
            "订单确认",
            $"您的订单 {notification.OrderId} 已创建,总金额: {notification.TotalAmount:C}");
    }
}

// 处理器 2:更新缓存
public class InvalidateOrderCacheHandler : INotificationHandler<OrderCreatedNotification>
{
    private readonly IDistributedCache _cache;

    public async Task Handle(OrderCreatedNotification notification, CancellationToken cancellationToken)
    {
        // 清除订单列表缓存
        await _cache.RemoveAsync("orders_list", cancellationToken);
        
        // 清除客户订单缓存
        await _cache.RemoveAsync($"customer_orders_{notification.CustomerEmail}", cancellationToken);
    }
}

// 处理器 3:写审计日志
public class AuditLogHandler : INotificationHandler<OrderCreatedNotification>
{
    private readonly IAuditService _auditService;

    public async Task Handle(OrderCreatedNotification notification, CancellationToken cancellationToken)
    {
        await _auditService.LogAsync(
            eventType: "OrderCreated",
            entityId: notification.OrderId,
            data: new { notification.TotalAmount, notification.CreatedAt });
    }
}

// 处理器 4:发送分析事件
public class TrackAnalyticsHandler : INotificationHandler<OrderCreatedNotification>
{
    private readonly IAnalyticsService _analyticsService;

    public async Task Handle(OrderCreatedNotification notification, CancellationToken cancellationToken)
    {
        await _analyticsService.TrackEventAsync("order_created", new
        {
            notification.OrderId,
            notification.TotalAmount,
            notification.CreatedAt
        });
    }
}

发布通知 ​

csharp
public class CreateOrderHandler : IRequestHandler<CreateOrderCommand, OrderResult>
{
    private readonly IOrderRepository _orderRepo;
    private readonly IMediator _mediator;

    public CreateOrderHandler(IOrderRepository orderRepo, IMediator mediator)
    {
        _orderRepo = orderRepo;
        _mediator = mediator;
    }

    public async Task<OrderResult> Handle(CreateOrderCommand request, CancellationToken cancellationToken)
    {
        // 1. 创建订单
        var order = new Order
        {
            ProductName = request.ProductName,
            Quantity = request.Quantity,
            TotalAmount = request.Price * request.Quantity
        };

        await _orderRepo.AddAsync(order, cancellationToken);

        // 2. 发布通知(所有处理器会被调用)
        await _mediator.Publish(new OrderCreatedNotification
        {
            OrderId = order.Id,
            CustomerEmail = request.CustomerEmail,
            TotalAmount = order.TotalAmount,
            CreatedAt = DateTime.UtcNow
        }, cancellationToken);

        return new OrderResult { OrderId = order.Id };
    }
}

⚡ 顺序执行 vs 并行执行 ​

默认行为:并行执行 ​

MediatR 默认并行执行所有通知处理器:

csharp
// 内部实现(简化版)
protected override async Task PublishCore(IEnumerable<NotificationHandlerExecutor> handlers, ...)
{
    // 并行执行所有处理器
    var tasks = handlers.Select(h => h.HandlerCallback(notification, cancellationToken));
    await Task.WhenAll(tasks);
}

优点:

  • ✅ 性能更好(尤其是 I/O 密集型操作)
  • ✅ 处理器之间互不影响

缺点:

  • ❌ 执行顺序不确定
  • ❌ 一个处理器异常会影响其他处理器

自定义:顺序执行 ​

如果需要顺序执行,可以自定义 Mediator:

csharp
public class SequentialMediator : Mediator
{
    private readonly ILogger<SequentialMediator> _logger;

    public SequentialMediator(IServiceProvider serviceFactory, ILogger<SequentialMediator> logger) 
        : base(serviceFactory)
    {
        _logger = logger;
    }

    protected override async Task PublishCore(
        IEnumerable<NotificationHandlerExecutor> handlerExecutors, 
        INotification notification, 
        CancellationToken cancellationToken)
    {
        // 顺序执行每个处理器
        foreach (var handlerExecutor in handlerExecutors)
        {
            try
            {
                _logger.LogInformation("执行处理器: {HandlerType}", 
                    handlerExecutor.HandlerInstance.GetType().Name);
                
                await handlerExecutor.HandlerCallback(notification, cancellationToken);
                
                _logger.LogInformation("处理器执行完成: {HandlerType}", 
                    handlerExecutor.HandlerInstance.GetType().Name);
            }
            catch (Exception ex)
            {
                _logger.LogError(ex, "处理器执行失败: {HandlerType}", 
                    handlerExecutor.HandlerInstance.GetType().Name);
                
                // 选项 1:继续执行其他处理器
                continue;
                
                // 选项 2:中断后续处理器
                // throw;
            }
        }
    }
}

注册自定义 Mediator:

csharp
builder.Services.AddScoped<IMediator, SequentialMediator>();

⚠️ 异常处理与容错 ​

问题:一个处理器异常会中断其他处理器 ​

csharp
// 默认行为
await _mediator.Publish(notification);

// 如果处理器 1 抛出异常,处理器 2、3、4 可能不会执行

解决方案 1:自定义 Mediator(推荐) ​

csharp
public class ResilientMediator : Mediator
{
    private readonly ILogger<ResilientMediator> _logger;

    public ResilientMediator(IServiceProvider serviceFactory, ILogger<ResilientMediator> logger) 
        : base(serviceFactory)
    {
        _logger = logger;
    }

    protected override async Task PublishCore(
        IEnumerable<NotificationHandlerExecutor> handlerExecutors, 
        INotification notification, 
        CancellationToken cancellationToken)
    {
        foreach (var handlerExecutor in handlerExecutors)
        {
            try
            {
                await handlerExecutor.HandlerCallback(notification, cancellationToken);
            }
            catch (Exception ex)
            {
                // 记录错误但继续执行其他处理器
                _logger.LogError(ex, 
                    "通知处理器失败: {HandlerType}. 通知类型: {NotificationType}",
                    handlerExecutor.HandlerInstance.GetType().Name,
                    notification.GetType().Name);
            }
        }
    }
}

解决方案 2:在处理器内部捕获异常 ​

csharp
public class SendEmailHandler : INotificationHandler<OrderCreatedNotification>
{
    public async Task Handle(OrderCreatedNotification notification, CancellationToken cancellationToken)
    {
        try
        {
            await _emailService.SendAsync(/* ... */);
        }
        catch (Exception ex)
        {
            // 记录错误但不抛出,避免影响其他处理器
            _logger.LogError(ex, "发送邮件失败,但不影响其他处理器");
        }
    }
}

🔥 领域事件与通知的结合 ​

什么是领域事件? ​

领域事件(Domain Events)是领域驱动设计(DDD)中的重要概念,表示领域中发生的重要事情。

实现领域事件 ​

csharp
// 领域事件基类
public abstract class DomainEvent : INotification
{
    public Guid EventId { get; } = Guid.NewGuid();
    public DateTime OccurredOn { get; } = DateTime.UtcNow;
}

// 具体领域事件
public class OrderCreatedEvent : DomainEvent
{
    public Guid OrderId { get; set; }
    public Guid CustomerId { get; set; }
    public decimal TotalAmount { get; set; }
}

public class OrderPaidEvent : DomainEvent
{
    public Guid OrderId { get; set; }
    public string PaymentMethod { get; set; }
    public DateTime PaidAt { get; set; }
}

public class OrderShippedEvent : DomainEvent
{
    public Guid OrderId { get; set; }
    public string TrackingNumber { get; set; }
    public DateTime ShippedAt { get; set; }
}

在聚合根中收集领域事件 ​

csharp
public class Order : Entity
{
    private readonly List<DomainEvent> _domainEvents = new();

    public IReadOnlyCollection<DomainEvent> DomainEvents => _domainEvents.AsReadOnly();

    public void AddDomainEvent(DomainEvent domainEvent)
    {
        _domainEvents.Add(domainEvent);
    }

    public void ClearDomainEvents()
    {
        _domainEvents.Clear();
    }

    // 业务方法
    public void Pay(string paymentMethod)
    {
        // 业务逻辑
        Status = OrderStatus.Paid;
        PaidAt = DateTime.UtcNow;

        // 添加领域事件
        AddDomainEvent(new OrderPaidEvent
        {
            OrderId = this.Id,
            PaymentMethod = paymentMethod,
            PaidAt = PaidAt
        });
    }
}

发布领域事件 ​

csharp
public class DomainEventPublisher
{
    private readonly IMediator _mediator;

    public DomainEventPublisher(IMediator mediator)
    {
        _mediator = mediator;
    }

    public async Task PublishEvents(Entity entity)
    {
        if (entity is Order order)
        {
            foreach (var domainEvent in order.DomainEvents)
            {
                await _mediator.Publish(domainEvent);
            }
            
            order.ClearDomainEvents();
        }
    }
}

🎯 最佳实践 ​

✅ 推荐做法 ​

csharp
// 1. 通知类保持简洁(只包含数据)
public class OrderCreatedNotification : INotification
{
    public Guid OrderId { get; set; }
    public decimal TotalAmount { get; set; }
}

// 2. 处理器单一职责
public class SendEmailHandler : INotificationHandler<OrderCreatedNotification>
{
    // 只负责发邮件
}

// 3. 使用容错 Mediator
builder.Services.AddScoped<IMediator, ResilientMediator>();

// 4. 异步处理
public async Task Handle(OrderCreatedNotification notification, CancellationToken ct)
{
    await _emailService.SendAsync(/* ... */, ct);
}

❌ 避免的做法 ​

csharp
// 1. 不要在通知中包含复杂逻辑
public class OrderCreatedNotification : INotification
{
    public void Process() // ❌ 错误
    {
        // 不要在通知中执行业务逻辑
    }
}

// 2. 不要依赖处理器的执行顺序
// 如果需要顺序,使用 SequentialMediator

// 3. 不要在处理器中执行耗时同步操作
public Task Handle(OrderCreatedNotification n, CancellationToken ct)
{
    Thread.Sleep(5000); // ❌ 阻塞线程
    return Task.CompletedTask;
}

📊 应用场景 ​

场景通知示例处理器
订单创建OrderCreatedNotification发邮件、更新缓存、写审计日志、发送分析事件
用户注册UserRegisteredNotification发欢迎邮件、初始化数据、送优惠券
支付成功PaymentSuccessNotification更新库存、生成发票、通知物流
文章发布ArticlePublishedNotification推送通知、更新索引、同步到 CDN

🎓 总结 ​

核心价值 ​

  1. 解耦:发布者和订阅者完全独立
  2. 扩展性:新增处理器无需修改代码
  3. 灵活性:支持并行和顺序执行
  4. 容错性:可自定义异常处理策略

与请求的区别 ​

特性请求通知
处理器数量1个N个
返回值有无
执行模式点对点发布/订阅
典型用途CQRS领域事件

💡 提示:通知机制是实现事件驱动架构和领域驱动设计的核心工具!

Released under the CC BY-SA 4.0 License.