通知(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个 | N个 |
| 返回值 | 有 | 无 |
| 执行模式 | 点对点 | 发布/订阅 |
| 典型用途 | CQRS | 领域事件 |
💡 提示:通知机制是实现事件驱动架构和领域驱动设计的核心工具!