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);
}
}
}关键发现:
- ✅ ServiceFactory 模式:使用委托而非直接依赖 DI 容器
- ✅ 包装器模式:通过
RequestHandlerWrapper封装处理逻辑 - ✅ 反射创建包装器:运行时动态创建泛型类型
- ✅ 虚拟方法:
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;
}
}关键机制:
- ✅ 延迟解析:Handler 在首次调用时才从容器解析
- ✅ 管道构建:反向遍历行为列表,形成嵌套委托链
- ✅ 装饰器模式:每个 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)));
}关键特性:
- ✅ 多播机制:返回所有注册的处理器
- ✅ 顺序执行:默认按注册顺序执行
- ✅ 可扩展:
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🎓 总结
核心机制
- ✅ ServiceFactory 委托:解耦 DI 容器
- ✅ 包装器模式:封装 Handler 和管道
- ✅ 装饰器链:反向构建行为管道
- ✅ 多播通知:顺序执行所有处理器
- ✅ 虚方法扩展:支持自定义行为
性能特征
- ⚡ 单次请求开销:~100ns(主要是反射和委托调用)
- 💾 内存分配:~128 字节/请求
- 🔄 无内置缓存:依赖 DI 容器生命周期
- 🚀 源生成器可消除反射开销
扩展点
- 🎯 重写
PublishCore:自定义发布策略 - 🎯 自定义
ServiceFactory:集成其他 DI 容器 - 🎯 源生成器:编译时代码生成
💡 提示: 理解源码有助于更好地使用和优化 MediatR!