与其他框架/库集成
📋 概述
MediatR 的设计理念使其能够与众多 .NET 生态系统中的优秀库无缝集成。本文将深入探讨如何将这些库与 MediatR 结合使用,构建更强大的应用程序。
学习目标
- ✅ 掌握与 Autofac 容器集成
- ✅ 学习 Scrutor 装饰器模式
- ✅ FluentValidation 深度集成
- ✅ Polly 实现重试/熔断
- ✅ Serilog 和 Application Insights 日志集成
🔧 与 Autofac 容器集成
为什么选择 Autofac?
Autofac 是一个功能强大的 IoC 容器,提供了比默认 DI 容器更多的特性:
- 属性注入
- 动态代理
- 模块化注册
- 生命周期作用域控制
- AOP 支持
安装 NuGet 包
bash
dotnet add package Autofac
dotnet add package Autofac.Extensions.DependencyInjection
dotnet add package MediatR配置 Autofac
csharp
// Program.cs
using Autofac;
using Autofac.Extensions.DependencyInjection;
using MediatR;
var builder = WebApplication.CreateBuilder(args);
// 使用 Autofac 替代默认容器
builder.Host.UseServiceProviderFactory(new AutofacServiceProviderFactory());
// 注册 MediatR
builder.Services.AddMediatR(cfg =>
{
cfg.RegisterServicesFromAssembly(typeof(Program).Assembly);
});
// Autofac 模块配置
builder.Host.ConfigureContainer<ContainerBuilder>(containerBuilder =>
{
// 注册 MediatR 服务到 Autofac
containerBuilder.RegisterModule(new MediatRModule());
// 注册其他服务
containerBuilder.RegisterType<OrderRepository>()
.As<IOrderRepository>()
.InstancePerLifetimeScope();
containerBuilder.RegisterType<CustomerService>()
.As<ICustomerService>()
.InstancePerLifetimeScope();
});
var app = builder.Build();
app.Run();
// MediatR 模块
public class MediatRModule : Module
{
protected override void Load(ContainerBuilder builder)
{
var assembly = Assembly.GetExecutingAssembly();
// 注册所有 Handler
builder.RegisterAssemblyTypes(assembly)
.AsClosedTypesOf(typeof(IRequestHandler<,>))
.AsImplementedInterfaces();
// 注册所有 Notification Handlers
builder.RegisterAssemblyTypes(assembly)
.AsClosedTypesOf(typeof(INotificationHandler<>))
.AsImplementedInterfaces();
// 注册 Pipeline Behaviors
builder.RegisterGeneric(typeof(LoggingBehavior<,>))
.As(typeof(IPipelineBehavior<,>))
.InstancePerDependency();
builder.RegisterGeneric(typeof(ValidationBehavior<,>))
.As(typeof(IPipelineBehavior<,>))
.InstancePerDependency();
}
}高级特性:属性注入
csharp
// 使用属性注入 IMediator
public class OrderController : ControllerBase
{
public IMediator Mediator { get; set; } = default!; // 属性注入
[HttpPost]
public async Task<ActionResult<int>> CreateOrder(CreateOrderCommand command)
{
var orderId = await Mediator.Send(command);
return CreatedAtAction(nameof(GetOrder), new { id = orderId }, orderId);
}
}
// Autofac 配置属性注入
builder.RegisterControllers(Assembly.GetExecutingAssembly())
.PropertiesAutowired();🎨 与 Scrutor 装饰器结合
什么是装饰器模式?
装饰器模式允许在不修改原有代码的情况下增强功能。Scrutor 提供了便捷的装饰器注册方式。
安装 NuGet 包
bash
dotnet add package Scrutor基本用法
csharp
// Program.cs
using Scrutor;
builder.Services.AddScoped<IOrderService, OrderService>();
// 使用 Scrutor 添加装饰器
builder.Services.Decorate<IOrderService, CachingOrderService>();
builder.Services.Decorate<IOrderService, LoggingOrderService>();
// 最终解析顺序:
// LoggingOrderService -> CachingOrderService -> OrderService与 MediatR 结合
csharp
// 定义接口
public interface IOrderProcessor
{
Task ProcessOrderAsync(int orderId);
}
// 核心实现
public class OrderProcessor : IOrderProcessor
{
private readonly IMediator _mediator;
public OrderProcessor(IMediator mediator)
{
_mediator = mediator;
}
public async Task ProcessOrderAsync(int orderId)
{
await _mediator.Send(new ProcessOrderCommand(orderId));
}
}
// 缓存装饰器
public class CachingOrderProcessor : IOrderProcessor
{
private readonly IOrderProcessor _inner;
private readonly IDistributedCache _cache;
public CachingOrderProcessor(IOrderProcessor inner, IDistributedCache cache)
{
_inner = inner;
_cache = cache;
}
public async Task ProcessOrderAsync(int orderId)
{
var cacheKey = $"order_processed:{orderId}";
if (await _cache.GetStringAsync(cacheKey) != null)
{
return; // 已处理,跳过
}
await _inner.ProcessOrderAsync(orderId);
await _cache.SetStringAsync(
cacheKey,
"true",
new DistributedCacheEntryOptions
{
AbsoluteExpirationRelativeToNow = TimeSpan.FromHours(1)
});
}
}
// 重试装饰器
public class RetryOrderProcessor : IOrderProcessor
{
private readonly IOrderProcessor _inner;
private readonly ILogger<RetryOrderProcessor> _logger;
public RetryOrderProcessor(IOrderProcessor inner, ILogger<RetryOrderProcessor> logger)
{
_inner = inner;
_logger = logger;
}
public async Task ProcessOrderAsync(int orderId)
{
var retryCount = 3;
for (int i = 0; i < retryCount; i++)
{
try
{
await _inner.ProcessOrderAsync(orderId);
return;
}
catch (Exception ex)
{
_logger.LogWarning(ex,
"订单处理失败,第 {RetryCount} 次重试", i + 1);
if (i == retryCount - 1)
throw;
await Task.Delay(TimeSpan.FromSeconds(Math.Pow(2, i)));
}
}
}
}
// 注册
builder.Services.AddScoped<IOrderProcessor, OrderProcessor>();
builder.Services.Decorate<IOrderProcessor, CachingOrderProcessor>();
builder.Services.Decorate<IOrderProcessor, RetryOrderProcessor>();
// 使用
public class OrdersController : ControllerBase
{
private readonly IOrderProcessor _orderProcessor;
public OrdersController(IOrderProcessor orderProcessor)
{
_orderProcessor = orderProcessor;
}
[HttpPost("{id}/process")]
public async Task<IActionResult> ProcessOrder(int id)
{
await _orderProcessor.ProcessOrderAsync(id);
return NoContent();
}
}✅ 与 FluentValidation 深度集成
自动验证管道行为
这是最常用的集成之一,实现请求的自动验证。
安装 NuGet 包
bash
dotnet add package FluentValidation
dotnet add package FluentValidation.DependencyInjectionExtensions定义验证器
csharp
// Commands/Orders/CreateOrder/CreateOrderValidator.cs
using FluentValidation;
public class CreateOrderValidator : AbstractValidator<CreateOrderCommand>
{
public CreateOrderValidator()
{
RuleFor(x => x.CustomerId)
.GreaterThan(0)
.WithMessage("客户ID必须大于0");
RuleFor(x => x.Items)
.NotEmpty()
.WithMessage("订单商品不能为空")
.Must(x => x.Count > 0)
.WithMessage("至少需要一个商品");
RuleForEach(x => x.Items).ChildRules(item =>
{
item.RuleFor(x => x.ProductId)
.GreaterThan(0)
.WithMessage("商品ID无效");
item.RuleFor(x => x.Quantity)
.GreaterThan(0)
.WithMessage("数量必须大于0");
item.RuleFor(x => x.UnitPrice)
.GreaterThanOrEqualTo(0)
.WithMessage("价格不能为负数");
});
RuleFor(x => x.ShippingAddress)
.NotNull()
.WithMessage("收货地址不能为空");
When(x => x.ShippingAddress != null, () =>
{
RuleFor(x => x.ShippingAddress.Province)
.NotEmpty()
.WithMessage("省份不能为空");
RuleFor(x => x.ShippingAddress.City)
.NotEmpty()
.WithMessage("城市不能为空");
});
}
}
// 复杂验证器(需要依赖注入)
public class CreateOrderAdvancedValidator : AbstractValidator<CreateOrderCommand>
{
private readonly ICustomerRepository _customerRepository;
private readonly IInventoryService _inventoryService;
public CreateOrderAdvancedValidator(
ICustomerRepository customerRepository,
IInventoryService inventoryService)
{
_customerRepository = customerRepository;
_inventoryService = inventoryService;
RuleFor(x => x.CustomerId)
.MustAsync(CustomerExists)
.WithMessage("客户不存在");
RuleFor(x => x.Items)
.MustAsync(AllItemsInStock)
.WithMessage("部分商品库存不足");
}
private async Task<bool> CustomerExists(int customerId, CancellationToken ct)
{
var customer = await _customerRepository.GetByIdAsync(customerId, ct);
return customer != null;
}
private async Task<bool> AllItemsInStock(
List<OrderItemDto> items,
CancellationToken ct)
{
foreach (var item in items)
{
var stock = await _inventoryService.GetStockAsync(item.ProductId, ct);
if (stock < item.Quantity)
return false;
}
return true;
}
}验证管道行为
csharp
// Behaviors/ValidationBehavior.cs
public class ValidationBehavior<TRequest, TResponse>
: IPipelineBehavior<TRequest, TResponse>
where TRequest : IRequest<TResponse>
{
private readonly IEnumerable<IValidator<TRequest>> _validators;
public ValidationBehavior(IEnumerable<IValidator<TRequest>> validators)
{
_validators = validators;
}
public async Task<TResponse> Handle(
TRequest request,
RequestHandlerDelegate<TResponse> next,
CancellationToken cancellationToken)
{
if (!_validators.Any())
return await next();
var context = new ValidationContext<TRequest>(request);
var validationResults = await Task.WhenAll(
_validators.Select(v => v.ValidateAsync(context, cancellationToken)));
var failures = validationResults
.SelectMany(r => r.Errors)
.Where(f => f != null)
.ToList();
if (failures.Any())
throw new ValidationException(failures);
return await next();
}
}注册验证器
csharp
// Program.cs
using FluentValidation;
// 自动扫描并注册验证器
builder.Services.AddValidatorsFromAssemblyContaining<Program>();
// 或者手动注册
builder.Services.AddScoped<IValidator<CreateOrderCommand>, CreateOrderValidator>();
builder.Services.AddScoped<IValidator<UpdateOrderCommand>, UpdateOrderValidator>();
// 注册验证行为
builder.Services.AddTransient(
typeof(IPipelineBehavior<,>),
typeof(ValidationBehavior<,>));全局异常处理
csharp
// Middleware/ValidationExceptionHandlerMiddleware.cs
public class ValidationExceptionHandlerMiddleware
{
private readonly RequestDelegate _next;
private readonly ILogger<ValidationExceptionHandlerMiddleware> _logger;
public ValidationExceptionHandlerMiddleware(
RequestDelegate next,
ILogger<ValidationExceptionHandlerMiddleware> logger)
{
_next = next;
_logger = logger;
}
public async Task InvokeAsync(HttpContext context)
{
try
{
await _next(context);
}
catch (ValidationException ex)
{
_logger.LogWarning(ex, "验证失败");
context.Response.StatusCode = StatusCodes.Status400BadRequest;
context.Response.ContentType = "application/json";
var errors = ex.Errors.ToDictionary(
e => e.PropertyName,
e => new[] { e.ErrorMessage });
var response = new
{
type = "https://tools.ietf.org/html/rfc7231#section-6.5.1",
title = "One or more validation errors occurred.",
status = 400,
errors = errors
};
await context.Response.WriteAsJsonAsync(response);
}
}
}
// 使用中间件
app.UseMiddleware<ValidationExceptionHandlerMiddleware>();🔄 与 Polly 实现重试/熔断
为什么需要 Polly?
Polly 是一个弹性和瞬态故障处理库,可以帮助处理:
- 重试(Retry)
- 断路器(Circuit Breaker)
- 超时(Timeout)
- 隔离(Bulkhead Isolation)
- 缓存(Cache)
- 回退(Fallback)
安装 NuGet 包
bash
dotnet add package Polly
dotnet add package Polly.Extensions.HttpPolly 管道行为
csharp
// Behaviors/PollyBehavior.cs
using Polly;
using Polly.Retry;
public class PollyBehavior<TRequest, TResponse>
: IPipelineBehavior<TRequest, TResponse>
where TRequest : IRequest<TResponse>
{
private readonly AsyncRetryPolicy _retryPolicy;
private readonly ICircuitBreakerPolicy _circuitBreakerPolicy;
private readonly ILogger<PollyBehavior<TRequest, TResponse>> _logger;
public PollyBehavior(ILogger<PollyBehavior<TRequest, TResponse>> logger)
{
_logger = logger;
// 重试策略:最多3次,指数退避
_retryPolicy = Policy
.Handle<Exception>()
.WaitAndRetryAsync(
retryCount: 3,
sleepDurationProvider: attempt =>
TimeSpan.FromSeconds(Math.Pow(2, attempt)),
onRetry: (exception, timeSpan, retryCount, context) =>
{
_logger.LogWarning(exception,
"请求 {RequestType} 失败,第 {RetryCount} 次重试",
typeof(TRequest).Name,
retryCount);
});
// 熔断策略:5次失败后熔断30秒
_circuitBreakerPolicy = Policy
.Handle<Exception>()
.CircuitBreakerAsync(
exceptionsAllowedBeforeBreaking: 5,
durationOfBreak: TimeSpan.FromSeconds(30),
onBreak: (exception, duration) =>
{
_logger.LogError(exception,
"熔断器打开,持续 {Duration} 秒",
duration.TotalSeconds);
},
onReset: () =>
{
_logger.LogInformation("熔断器重置");
});
}
public async Task<TResponse> Handle(
TRequest request,
RequestHandlerDelegate<TResponse> next,
CancellationToken cancellationToken)
{
return await _circuitBreakerPolicy.ExecuteAsync(() =>
_retryPolicy.ExecuteAsync(async () =>
{
return await next();
}));
}
}
// 注册
builder.Services.AddTransient(
typeof(IPipelineBehavior<,>),
typeof(PollyBehavior<,>));针对特定请求的策略
csharp
// 定义标记接口
public interface IResilientRequest { }
// 命令实现标记接口
public record ProcessPaymentCommand : IRequest<bool>, IResilientRequest
{
public int OrderId { get; init; }
public decimal Amount { get; init; }
}
// 弹性行为(只处理标记的请求)
public class ResilienceBehavior<TRequest, TResponse>
: IPipelineBehavior<TRequest, TResponse>
where TRequest : IRequest<TResponse>, IResilientRequest
{
private readonly IAsyncPolicy<TResponse> _policy;
public ResilienceBehavior()
{
_policy = Policy<TResponse>
.Handle<Exception>()
.OrResult(result => Equals(result, default(TResponse)))
.WaitAndRetryAsync(
retryCount: 3,
sleepDurationProvider: attempt =>
TimeSpan.FromSeconds(Math.Pow(2, attempt)));
}
public async Task<TResponse> Handle(
TRequest request,
RequestHandlerDelegate<TResponse> next,
CancellationToken cancellationToken)
{
return await _policy.ExecuteAsync(() => next());
}
}Timeout 行为
csharp
// Behaviors/TimeoutBehavior.cs
public class TimeoutBehavior<TRequest, TResponse>
: IPipelineBehavior<TRequest, TResponse>
where TRequest : IRequest<TResponse>
{
private readonly ILogger<TimeoutBehavior<TRequest, TResponse>> _logger;
public async Task<TResponse> Handle(
TRequest request,
RequestHandlerDelegate<TResponse> next,
CancellationToken cancellationToken)
{
using var cts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
cts.CancelAfter(TimeSpan.FromSeconds(30)); // 30秒超时
try
{
return await next();
}
catch (OperationCanceledException) when (cts.IsCancellationRequested)
{
_logger.LogError("请求 {RequestType} 超时", typeof(TRequest).Name);
throw new TimeoutException(
$"请求 {typeof(TRequest).Name} 执行超时(30秒)");
}
}
}📝 与 Serilog / Application Insights 集成
Serilog 结构化日志
安装 NuGet 包
bash
dotnet add package Serilog.AspNetCore
dotnet add package Serilog.Sinks.Console
dotnet add package Serilog.Sinks.File
dotnet add package Serilog.Enrichers.Environment配置 Serilog
csharp
// Program.cs
using Serilog;
Log.Logger = new LoggerConfiguration()
.MinimumLevel.Information()
.MinimumLevel.Override("Microsoft", LogEventLevel.Warning)
.MinimumLevel.Override("Microsoft.Hosting.Lifetime", LogEventLevel.Information)
.Enrich.FromLogContext()
.Enrich.WithEnvironmentName()
.Enrich.WithMachineName()
.WriteTo.Console(
outputTemplate: "[{Timestamp:HH:mm:ss} {Level:u3}] {Message:lj} {Properties:j}{NewLine}{Exception}")
.WriteTo.File(
path: "logs/log-.txt",
rollingInterval: RollingInterval.Day,
retainedFileCountLimit: 30,
outputTemplate: "{Timestamp:yyyy-MM-dd HH:mm:ss.fff zzz} [{Level:u3}] {Message:lj} {Properties:j}{NewLine}{Exception}")
.CreateLogger();
try
{
Log.Information("应用程序启动");
var builder = WebApplication.CreateBuilder(args);
builder.Host.UseSerilog();
// ... 其他配置
var app = builder.Build();
app.Run();
}
catch (Exception ex)
{
Log.Fatal(ex, "应用程序启动失败");
}
finally
{
Log.CloseAndFlush();
}增强的日志行为
csharp
// Behaviors/SerilogLoggingBehavior.cs
using Serilog.Context;
public class SerilogLoggingBehavior<TRequest, TResponse>
: IPipelineBehavior<TRequest, TResponse>
where TRequest : IRequest<TResponse>
{
private readonly ILogger<SerilogLoggingBehavior<TRequest, TResponse>> _logger;
public async Task<TResponse> Handle(
TRequest request,
RequestHandlerDelegate<TResponse> next,
CancellationToken cancellationToken)
{
var requestType = typeof(TRequest).Name;
var requestId = Guid.NewGuid();
// enrich log context
using (LogContext.PushProperty("RequestId", requestId))
using (LogContext.PushProperty("RequestType", requestType))
{
_logger.LogInformation(
"开始处理请求 [{RequestId}] {RequestType}: {@Request}",
requestId,
requestType,
request);
var stopwatch = Stopwatch.StartNew();
try
{
var response = await next();
stopwatch.Stop();
_logger.LogInformation(
"请求完成 [{RequestId}] {RequestType} - 耗时: {ElapsedMs}ms",
requestId,
requestType,
stopwatch.ElapsedMilliseconds);
return response;
}
catch (Exception ex)
{
stopwatch.Stop();
_logger.LogError(ex,
"请求失败 [{RequestId}] {RequestType} - 耗时: {ElapsedMs}ms",
requestId,
requestType,
stopwatch.ElapsedMilliseconds);
throw;
}
}
}
}Azure Application Insights 集成
安装 NuGet 包
bash
dotnet add package Microsoft.ApplicationInsights.AspNetCore配置 Application Insights
csharp
// Program.cs
builder.Services.AddApplicationInsightsTelemetry(options =>
{
options.ConnectionString = builder.Configuration["APPLICATIONINSIGHTS_CONNECTION_STRING"];
});
// 自定义遥测初始化器
builder.Services.AddSingleton<ITelemetryInitializer, MediatRTelemetryInitializer>();MediatR 遥测初始化器
csharp
// Diagnostics/MediatRTelemetryInitializer.cs
using Microsoft.ApplicationInsights.Channel;
using Microsoft.ApplicationInsights.DataContracts;
using Microsoft.ApplicationInsights.Extensibility;
public class MediatRTelemetryInitializer : ITelemetryInitializer
{
private readonly IHttpContextAccessor _httpContextAccessor;
public MediatRTelemetryInitializer(IHttpContextAccessor httpContextAccessor)
{
_httpContextAccessor = httpContextAccessor;
}
public void Initialize(ITelemetry telemetry)
{
if (telemetry is not RequestTelemetry requestTelemetry)
return;
// 添加 MediatR 相关维度
var httpContext = _httpContextAccessor.HttpContext;
if (httpContext?.Items.ContainsKey("MediatR_RequestType") == true)
{
var requestType = httpContext.Items["MediatR_RequestType"]?.ToString();
telemetry.Properties["MediatR.RequestType"] = requestType;
}
}
}
// 在管道行为中设置
public class AppInsightsBehavior<TRequest, TResponse>
: IPipelineBehavior<TRequest, TResponse>
where TRequest : IRequest<TResponse>
{
private readonly TelemetryClient _telemetryClient;
private readonly IHttpContextAccessor _httpContextAccessor;
public async Task<TResponse> Handle(
TRequest request,
RequestHandlerDelegate<TResponse> next,
CancellationToken cancellationToken)
{
var operation = _telemetryClient.StartOperation<DependencyTelemetry>(
$"MediatR {typeof(TRequest).Name}");
operation.Telemetry.Type = "MediatR";
operation.Telemetry.Data = JsonSerializer.Serialize(request);
// 存储到 HttpContext
_httpContextAccessor.HttpContext?.Items.TryAdd(
"MediatR_RequestType",
typeof(TRequest).Name);
try
{
var result = await next();
operation.Telemetry.Success = true;
return result;
}
catch (Exception ex)
{
operation.Telemetry.Success = false;
_telemetryClient.TrackException(ex);
throw;
}
finally
{
_telemetryClient.StopOperation(operation);
}
}
}🎯 完整集成示例
综合配置
csharp
// Program.cs
var builder = WebApplication.CreateBuilder(args);
// 1. Serilog
builder.Host.UseSerilog((context, config) =>
{
config.ReadFrom.Configuration(context.Configuration);
});
// 2. MediatR
builder.Services.AddMediatR(cfg =>
{
cfg.RegisterServicesFromAssembly(typeof(Program).Assembly);
});
// 3. FluentValidation
builder.Services.AddValidatorsFromAssemblyContaining<Program>();
// 4. Pipeline Behaviors
builder.Services.AddTransient(
typeof(IPipelineBehavior<,>),
typeof(SerilogLoggingBehavior<,>));
builder.Services.AddTransient(
typeof(IPipelineBehavior<,>),
typeof(ValidationBehavior<,>));
builder.Services.AddTransient(
typeof(IPipelineBehavior<,>),
typeof(PollyBehavior<,>));
// 5. Application Insights
builder.Services.AddApplicationInsightsTelemetry();
// 6. Scrutor Decorators
builder.Services.AddScoped<IOrderService, OrderService>();
builder.Services.Decorate<IOrderService, CachingOrderService>();
builder.Services.Decorate<IOrderService, LoggingOrderService>();
var app = builder.Build();
// 7. Exception Handling
app.UseMiddleware<ValidationExceptionHandlerMiddleware>();
app.UseSerilogRequestLogging();
app.Run();📊 对比总结
| 集成库 | 用途 | 复杂度 | 推荐度 |
|---|---|---|---|
| Autofac | 高级 DI 容器 | ⭐⭐⭐ | ⭐⭐⭐⭐ |
| Scrutor | 装饰器模式 | ⭐⭐ | ⭐⭐⭐⭐⭐ |
| FluentValidation | 自动验证 | ⭐⭐ | ⭐⭐⭐⭐⭐ |
| Polly | 弹性与容错 | ⭐⭐⭐ | ⭐⭐⭐⭐⭐ |
| Serilog | 结构化日志 | ⭐⭐ | ⭐⭐⭐⭐⭐ |
| App Insights | 监控与遥测 | ⭐⭐⭐ | ⭐⭐⭐⭐ |
🎓 最佳实践
✅ 推荐做法
- FluentValidation 必配:几乎所有项目都应该集成
- Serilog 优先:结构化日志便于排查问题
- Polly 用于外部调用:数据库、API 调用等
- Scrutor 简化装饰器:比手动实现更简洁
- Autofac 按需使用:复杂场景才需要
❌ 避免做法
- 不要过度集成:简单项目不需要所有库
- 避免重复功能:如多个日志库
- 注意性能影响:每个 Behavior 都有开销
- 不要忽略配置:正确的配置很重要