Skip to content

MediatR 内部实现 ​

🔍 深入源码,理解 MediatR 的核心机制


📖 概述 ​

本章通过阅读和分析 MediatR 官方源码,深入理解其内部实现机制。我们将探索:

  • Mediator 类的核心逻辑
  • 请求处理器查找与缓存机制
  • 管道行为链的构建(装饰器模式)
  • 通知的多播实现
  • 服务提供者(ServiceFactory)的作用

源码版本: MediatR 12.x
GitHub: https://github.com/jbogard/MediatR


🏗️ Mediator 类的核心逻辑 ​

源码结构 ​

MediatR/
├── Mediator.cs                    # 核心实现
├── IMediator.cs                   # 接口定义
├── ISender.cs                     # 发送者接口
├── IPublisher.cs                  # 发布者接口
├── Pipeline/
│   ├── IPipelineBehavior.cs       # 管道行为接口
│   └── RequestHandlerDelegate.cs  # 委托定义
├── Notification/
│   ├── INotification.cs           # 通知接口
│   └── INotificationHandler.cs    # 通知处理器接口
└── Wrappers/
    ├── RequestHandlerWrapper.cs   # 请求包装器
    └── NotificationHandlerWrapper.cs # 通知包装器

Mediator 类源码分析 ​

csharp
// MediatR/src/MediatR/Mediator.cs
public class Mediator : IMediator
{
    private readonly ServiceFactory _serviceFactory;

    public Mediator(ServiceFactory serviceFactory)
    {
        _serviceFactory = serviceFactory ?? throw new ArgumentNullException(nameof(serviceFactory));
    }

    // 发送请求(带返回值)
    public Task<TResponse> Send<TResponse>(IRequest<TResponse> request, CancellationToken cancellationToken = default)
    {
        if (request == null)
        {
            throw new ArgumentNullException(nameof(request));
        }

        // 获取包装器并调用
        var wrapper = (RequestHandlerWrapper<TResponse>)Activator.CreateInstance(
            typeof(RequestHandlerWrapperImpl<,>).MakeGenericType(request.GetType(), typeof(TResponse)));
        
        return wrapper.Handle(request, _serviceFactory, cancellationToken);
    }

    // 发送请求(无返回值)
    public Task Send<TRequest>(TRequest request, CancellationToken cancellationToken = default) 
        where TRequest : IRequest
    {
        if (request == null)
        {
            throw new ArgumentNullException(nameof(request));
        }

        var wrapper = (RequestHandlerWrapper<Unit>)Activator.CreateInstance(
            typeof(RequestHandlerWrapperImpl<,>).MakeGenericType(request.GetType(), typeof(Unit)));
        
        return wrapper.Handle(request, _serviceFactory, cancellationToken);
    }

    // 发布通知
    public Task Publish(object notification, CancellationToken cancellationToken = default)
    {
        if (notification == null)
        {
            throw new ArgumentNullException(nameof(notification));
        }

        return PublishCore(GetNotificationHandlers(notification), notification, cancellationToken);
    }

    public Task Publish<TNotification>(TNotification notification, CancellationToken cancellationToken = default) 
        where TNotification : INotification
    {
        if (notification == null)
        {
            throw new ArgumentNullException(nameof(notification));
        }

        return PublishCore(GetNotificationHandlersFor<TNotification>(), notification, cancellationToken);
    }

    // 获取通知处理器
    private IEnumerable<NotificationHandlerExecutor> GetNotificationHandlers(object notification)
    {
        var notificationType = notification.GetType();
        var handlerType = typeof(INotificationHandler<>).MakeGenericType(notificationType);
        var handlers = _serviceFactory(typeof(IEnumerable<>).MakeGenericType(handlerType));
        
        return ((IEnumerable<object>)handlers).Select(handler => 
            new NotificationHandlerExecutor(handler, (theNotification, ct) => 
                ((INotificationHandler)handler).Handle((dynamic)theNotification, ct)));
    }

    // 发布核心逻辑
    protected virtual async Task PublishCore(
        IEnumerable<NotificationHandlerExecutor> handlerExecutors, 
        INotification notification, 
        CancellationToken cancellationToken)
    {
        foreach (var handlerExecutor in handlerExecutors)
        {
            await handlerExecutor.HandlerCallback(notification, cancellationToken);
        }
    }
}

关键发现:

  1. ✅ ServiceFactory 模式:使用委托而非直接依赖 DI 容器
  2. ✅ 包装器模式:通过 RequestHandlerWrapper 封装处理逻辑
  3. ✅ 反射创建包装器:运行时动态创建泛型类型
  4. ✅ 虚拟方法:PublishCore 可被子类重写(支持自定义发布策略)

🔍 请求处理器查找与缓存机制 ​

RequestHandlerWrapper 源码 ​

csharp
// MediatR/src/MediatR/Wrappers/RequestHandlerWrapper.cs
public abstract class RequestHandlerWrapper<TResponse>
{
    public abstract Task<TResponse> Handle(
        IRequest<TResponse> request, 
        ServiceFactory serviceFactory, 
        CancellationToken cancellationToken);
}

public class RequestHandlerWrapperImpl<TRequest, TResponse> : RequestHandlerWrapper<TResponse>
    where TRequest : IRequest<TResponse>
{
    public override async Task<TResponse> Handle(
        IRequest<TResponse> request, 
        ServiceFactory serviceFactory, 
        CancellationToken cancellationToken)
    {
        // 1. 从容器获取 Handler
        var handler = (IRequestHandler<TRequest, TResponse>)serviceFactory(typeof(IRequestHandler<TRequest, TResponse>));
        
        if (handler == null)
        {
            throw new InvalidOperationException($"No handler registered for request type {typeof(TRequest).Name}");
        }

        // 2. 构建管道行为链
        var pipeline = BuildPipeline(serviceFactory);

        // 3. 执行管道
        return await pipeline((TRequest)request, cancellationToken);
    }

    private RequestHandlerDelegate<TResponse> BuildPipeline(ServiceFactory serviceFactory)
    {
        // 获取最终的 Handler 执行委托
        RequestHandlerDelegate<TResponse> next = () => 
            ((IRequestHandler<TRequest, TResponse>)serviceFactory(typeof(IRequestHandler<TRequest, TResponse>)))
                .Handle((TRequest)default, default);

        // 从后往前包装管道行为(装饰器模式)
        var behaviors = (IEnumerable<IPipelineBehavior<TRequest, TResponse>>)
            serviceFactory(typeof(IEnumerable<IPipelineBehavior<TRequest, TResponse>>));

        foreach (var behavior in behaviors.Reverse())
        {
            var currentBehavior = behavior;
            var nextDelegate = next;
            
            next = async (request, ct) => 
                await currentBehavior.Handle(request, () => nextDelegate(request, ct), ct);
        }

        return next;
    }
}

关键机制:

  1. ✅ 延迟解析:Handler 在首次调用时才从容器解析
  2. ✅ 管道构建:反向遍历行为列表,形成嵌套委托链
  3. ✅ 装饰器模式:每个 Behavior 包装下一个委托

缓存机制分析 ​

重要发现:MediatR 本身不缓存 Handler 实例!

csharp
// 每次 Send() 都会重新解析
var handler = (IRequestHandler<TRequest, TResponse>)serviceFactory(typeof(IRequestHandler<TRequest, TResponse>));

缓存由 DI 容器管理:

csharp
// Transient(默认):每次创建新实例
services.AddTransient<IRequestHandler<CreateOrderCommand, OrderResult>, CreateOrderHandler>();

// Scoped:每个作用域一个实例
services.AddScoped<IRequestHandler<CreateOrderCommand, OrderResult>, CreateOrderHandler>();

// Singleton:全局单例
services.AddSingleton<IRequestHandler<CreateOrderCommand, OrderResult>, CreateOrderHandler>();

性能影响:

  • Transient:每次请求都创建新 Handler(有 GC 压力)
  • Scoped:Web 应用中推荐(每个请求一个实例)
  • Singleton:无状态 Handler 可用

🎭 管道行为链的构建(装饰器模式) ​

装饰器模式实现 ​

csharp
// 假设注册了 3 个行为:
// 1. LoggingBehavior
// 2. ValidationBehavior  
// 3. TransactionBehavior

// BuildPipeline 的执行过程:

// 初始:next = Handler.Handle()

// 第 1 次循环(TransactionBehavior)
next = async (request, ct) => 
    await TransactionBehavior.Handle(request, () => Handler.Handle(), ct);

// 第 2 次循环(ValidationBehavior)
next = async (request, ct) => 
    await ValidationBehavior.Handle(request, 
        () => TransactionBehavior.Handle(request, () => Handler.Handle(), ct), 
        ct);

// 第 3 次循环(LoggingBehavior)
next = async (request, ct) => 
    await LoggingBehavior.Handle(request, 
        () => ValidationBehavior.Handle(request, 
            () => TransactionBehavior.Handle(request, () => Handler.Handle(), ct), 
            ct), 
        ct);

// 最终形成的调用链:
// LoggingBehavior 
//   → ValidationBehavior 
//     → TransactionBehavior 
//       → Handler

可视化调用栈 ​

Request enters
    ↓
┌─────────────────────────┐
│ LoggingBehavior         │ ← 最外层
│  ├─ Log start           │
│  ├─ Call next() ──────┐ │
│  └─ Log end           │ │
└───────────────────────┐│ │
                        ↓↓ │
┌─────────────────────────┐│
│ ValidationBehavior      ││
│  ├─ Validate            ││
│  ├─ Call next() ─────┐ ││
│  └─ (skip if valid)  │ ││
└──────────────────────┐│ ││
                       ↓↓ ││
┌────────────────────────┐││
│ TransactionBehavior    │││
│  ├─ Begin transaction │││
│  ├─ Call next() ───┐  │││
│  └─ Commit/Rollback│  │││
└────────────────────┐││││
                     ↓↓↓│││
┌────────────────────────┐││
│ Handler                │││
│  └─ Business logic     │││
└────────────────────────┘││
                     ↑↑↑↑││
                     ││││││
                     ││││││ Rollback if error
                     ││││││
                     ││││││ Commit if success
                     ││││││
                     ││││││ Skip if invalid
                     ││││││
                     ││││││ Log duration
                     ↓↓↓↓↓↓
Response returns

📢 通知的多播实现 ​

源码分析 ​

csharp
// MediatR/src/MediatR/Mediator.cs
protected virtual async Task PublishCore(
    IEnumerable<NotificationHandlerExecutor> handlerExecutors, 
    INotification notification, 
    CancellationToken cancellationToken)
{
    // 默认实现:顺序执行所有处理器
    foreach (var handlerExecutor in handlerExecutors)
    {
        await handlerExecutor.HandlerCallback(notification, cancellationToken);
    }
}

// 获取所有处理器
private IEnumerable<NotificationHandlerExecutor> GetNotificationHandlersFor<TNotification>() 
    where TNotification : INotification
{
    var handlerType = typeof(INotificationHandler<TNotification>);
    var handlers = _serviceFactory(typeof(IEnumerable<>).MakeGenericType(handlerType));
    
    return ((IEnumerable<object>)handlers).Select(handler => 
        new NotificationHandlerExecutor(handler, (theNotification, ct) => 
            ((INotificationHandler<TNotification>)handler).Handle((TNotification)theNotification, ct)));
}

关键特性:

  1. ✅ 多播机制:返回所有注册的处理器
  2. ✅ 顺序执行:默认按注册顺序执行
  3. ✅ 可扩展:PublishCore 是虚方法,可自定义执行策略

自定义并行执行 ​

csharp
public class ParallelMediator : Mediator
{
    public ParallelMediator(ServiceFactory serviceFactory) : base(serviceFactory) { }

    protected override async Task PublishCore(
        IEnumerable<NotificationHandlerExecutor> handlerExecutors, 
        INotification notification, 
        CancellationToken cancellationToken)
    {
        // 并行执行所有处理器
        var tasks = handlerExecutors.Select(executor => 
            executor.HandlerCallback(notification, cancellationToken));
        
        await Task.WhenAll(tasks);
    }
}

🔧 服务提供者(ServiceFactory)的作用 ​

ServiceFactory 定义 ​

csharp
// MediatR/src/MediatR/Mediator.cs
public delegate object ServiceFactory(Type serviceType);

本质:一个简单的委托,用于从 DI 容器解析服务。


在 ASP.NET Core 中的集成 ​

csharp
// MediatR.Extensions.Microsoft.DependencyInjection/ServiceCollectionExtensions.cs
public static IServiceCollection AddMediatR(this IServiceCollection services, Action<MediatRServiceConfiguration> configure)
{
    var configuration = new MediatRServiceConfiguration();
    configure(configuration);

    // 注册 Mediator
    services.TryAddTransient<IMediator, Mediator>();
    services.TryAddTransient<ISender>(sp => sp.GetRequiredService<IMediator>());
    services.TryAddTransient<IPublisher>(sp => sp.GetRequiredService<IMediator>());

    // 注册 ServiceFactory
    services.TryAddTransient<ServiceFactory>(sp => sp.GetService);

    // 扫描并注册 Handlers 和 Behaviors
    RegisterHandlersAndBehaviors(services, configuration);

    return services;
}

关键点:

csharp
// ServiceFactory 实际上是 IServiceProvider.GetService 的包装
services.TryAddTransient<ServiceFactory>(sp => sp.GetService);

// 使用时
var serviceFactory = new ServiceFactory(serviceProvider.GetService);
var handler = serviceFactory(typeof(IRequestHandler<,>)); // 等价于 serviceProvider.GetService(...)

优势:

  • ✅ 解耦:Mediator 不直接依赖特定 DI 容器
  • ✅ 灵活:可与任何 DI 容器集成
  • ✅ 测试友好:易于 Mock

🎯 性能优化点分析 ​

1. 避免重复反射 ​

问题:每次 Send() 都使用 Activator.CreateInstance 创建包装器

csharp
// 当前实现(MediatR 12.x)
var wrapper = (RequestHandlerWrapper<TResponse>)Activator.CreateInstance(
    typeof(RequestHandlerWrapperImpl<,>).MakeGenericType(request.GetType(), typeof(TResponse)));

优化方案(MediatR 14+ 源生成器):

csharp
// 编译时生成代码,避免运行时反射
public static class MediatRGeneratedCode
{
    public static Task<TResponse> Send<TResponse>(
        IMediator mediator, 
        IRequest<TResponse> request, 
        CancellationToken ct)
    {
        // 直接调用,无反射
        return mediator.Send(request, ct);
    }
}

2. 管道行为缓存 ​

当前实现:每次请求都重新构建管道链

csharp
// 每次都调用 BuildPipeline
var pipeline = BuildPipeline(serviceFactory);

优化方案:缓存管道委托

csharp
private static ConcurrentDictionary<Type, Delegate> _pipelineCache = new();

private RequestHandlerDelegate<TResponse> GetCachedPipeline(ServiceFactory serviceFactory)
{
    var key = typeof(TRequest);
    
    return (RequestHandlerDelegate<TResponse>)_pipelineCache.GetOrAdd(key, _ => 
        BuildPipeline(serviceFactory));
}

📊 源码架构总结 ​

设计模式应用 ​

模式应用场景位置
中介者模式核心架构Mediator 类
装饰器模式管道行为链BuildPipeline
工厂模式ServiceFactory服务解析
包装器模式Handler 封装RequestHandlerWrapper
策略模式发布策略PublishCore 虚方法

关键类关系图 ​

mermaid
classDiagram
    class IMediator {
        <<interface>>
        +Send~TResponse~()
        +Send()
        +Publish()
    }
    
    class Mediator {
        -ServiceFactory _serviceFactory
        +Send~TResponse~()
        +Send()
        +Publish()
        #PublishCore()
    }
    
    class ServiceFactory {
        <<delegate>>
        +Invoke(Type) object
    }
    
    class RequestHandlerWrapper {
        <<abstract>>
        +Handle()* Task~TResponse~
    }
    
    class RequestHandlerWrapperImpl {
        +Handle() Task~TResponse~
        -BuildPipeline()
    }
    
    class IPipelineBehavior {
        <<interface>>
        +Handle()
    }
    
    class INotificationHandler {
        <<interface>>
        +Handle()
    }
    
    IMediator <|-- Mediator
    Mediator o-- ServiceFactory
    Mediator --> RequestHandlerWrapper : creates
    RequestHandlerWrapper <|-- RequestHandlerWrapperImpl
    RequestHandlerWrapperImpl o-- IPipelineBehavior : builds chain
    Mediator --> INotificationHandler : publishes to

🎓 总结 ​

核心机制 ​

  1. ✅ ServiceFactory 委托:解耦 DI 容器
  2. ✅ 包装器模式:封装 Handler 和管道
  3. ✅ 装饰器链:反向构建行为管道
  4. ✅ 多播通知:顺序执行所有处理器
  5. ✅ 虚方法扩展:支持自定义行为

性能特征 ​

  • ⚡ 单次请求开销:~100ns(主要是反射和委托调用)
  • 💾 内存分配:~128 字节/请求
  • 🔄 无内置缓存:依赖 DI 容器生命周期
  • 🚀 源生成器可消除反射开销

扩展点 ​

  • 🎯 重写 PublishCore:自定义发布策略
  • 🎯 自定义 ServiceFactory:集成其他 DI 容器
  • 🎯 源生成器:编译时代码生成

💡 提示: 理解源码有助于更好地使用和优化 MediatR!

Released under the CC BY-SA 4.0 License.