管道行为(Pipeline Behaviors)
🔧 MediatR 最强大的特性 - 实现 AOP(面向切面编程)
📖 概述
管道行为(Pipeline Behaviors) 是 MediatR 中最强大和最灵活的特性之一。它允许你在请求处理前后插入自定义逻辑,实现横切关注点(Cross-Cutting Concerns)的分离。
什么是横切关注点?
横切关注点是那些影响多个模块的功能,例如:
- 📝 日志记录
- ✅ 参数验证
- 💾 事务管理
- 💾 缓存
- 📊 性能监控
- 🔒 安全检查
- 🔄 重试机制
传统方式中,这些逻辑会分散在各个方法中,导致代码重复和难以维护。管道行为通过装饰器模式将这些逻辑集中管理。
🎯 IPipelineBehavior<TRequest, TResponse> 原理
接口定义
csharp
public interface IPipelineBehavior<in TRequest, TResponse> where TRequest : notnull
{
Task<TResponse> Handle(
TRequest request,
RequestHandlerDelegate<TResponse> next,
CancellationToken cancellationToken);
}核心概念
| 参数 | 类型 | 说明 |
|---|---|---|
request | TRequest | 当前请求对象 |
next | RequestHandlerDelegate<TResponse> | 委托,指向下一个处理器(可能是另一个 Behavior 或最终的 Handler) |
cancellationToken | CancellationToken | 取消令牌 |
| 返回值 | Task<TResponse> | 处理结果 |
执行流程
mermaid
graph LR
A[Request] --> B[LoggingBehavior]
B --> C[ValidationBehavior]
C --> D[TransactionBehavior]
D --> E[CachingBehavior]
E --> F[Handler]
F --> E
E --> D
D --> C
C --> B
B --> G[Response]关键点:
- 请求依次通过每个 Behavior(外层 → 内层)
- 到达最终的 Handler
- 响应依次返回(内层 → 外层)
📝 日志行为(Logging Behavior)
实现
csharp
public class LoggingBehavior<TRequest, TResponse> : IPipelineBehavior<TRequest, TResponse>
where TRequest : IRequest<TResponse>
{
private readonly ILogger<LoggingBehavior<TRequest, TResponse>> _logger;
public LoggingBehavior(ILogger<LoggingBehavior<TRequest, TResponse>> logger)
{
_logger = logger;
}
public async Task<TResponse> Handle(
TRequest request,
RequestHandlerDelegate<TResponse> next,
CancellationToken cancellationToken)
{
var requestName = typeof(TRequest).Name;
var requestId = Guid.NewGuid();
_logger.LogInformation(
"[START] 处理请求 {RequestId}: {RequestName} {@Request}",
requestId, requestName, request);
var stopwatch = Stopwatch.StartNew();
try
{
var response = await next();
stopwatch.Stop();
_logger.LogInformation(
"[END] 请求 {RequestId}: {RequestName} 处理完成,耗时 {ElapsedMilliseconds}ms",
requestId, requestName, stopwatch.ElapsedMilliseconds);
return response;
}
catch (Exception ex)
{
stopwatch.Stop();
_logger.LogError(
ex,
"[ERROR] 请求 {RequestId}: {RequestName} 处理失败,耗时 {ElapsedMilliseconds}ms",
requestId, requestName, stopwatch.ElapsedMilliseconds);
throw; // 重新抛出异常
}
}
}注册
csharp
builder.Services.AddTransient(typeof(IPipelineBehavior<,>), typeof(LoggingBehavior<,>));输出示例
[START] 处理请求 a3f8b2c1: CreateOrderCommand {"ProductName":"iPhone","Quantity":2,"Price":7999}
[END] 请求 a3f8b2c1: CreateOrderCommand 处理完成,耗时 125ms✅ 验证行为(Validation Behavior)
安装 FluentValidation
bash
dotnet add package FluentValidation.AspNetCore
dotnet add package FluentValidation.DependencyInjectionExtensions定义验证器
csharp
public class CreateOrderCommandValidator : AbstractValidator<CreateOrderCommand>
{
public CreateOrderCommandValidator()
{
RuleFor(x => x.ProductName)
.NotEmpty().WithMessage("产品名称不能为空")
.MaximumLength(100).WithMessage("产品名称不能超过100个字符");
RuleFor(x => x.Quantity)
.GreaterThan(0).WithMessage("数量必须大于0")
.LessThanOrEqualTo(1000).WithMessage("数量不能超过1000");
RuleFor(x => x.Price)
.GreaterThan(0).WithMessage("价格必须大于0");
}
}实现验证行为
csharp
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.Count != 0)
{
throw new ValidationException(failures);
}
return await next();
}
}注册验证器和行为
csharp
// 注册验证器(自动扫描)
builder.Services.AddValidatorsFromAssembly(typeof(Program).Assembly);
// 注册验证行为
builder.Services.AddTransient(typeof(IPipelineBehavior<,>), typeof(ValidationBehavior<,>));错误响应示例
json
{
"errors": [
{
"propertyName": "ProductName",
"errorMessage": "产品名称不能为空"
},
{
"propertyName": "Quantity",
"errorMessage": "数量必须大于0"
}
]
}💾 事务行为(Transaction / Unit of Work Behavior)
实现
csharp
public class TransactionBehavior<TRequest, TResponse> : IPipelineBehavior<TRequest, TResponse>
where TRequest : IRequest<TResponse>
{
private readonly DbContext _dbContext;
private readonly ILogger<TransactionBehavior<TRequest, TResponse>> _logger;
public TransactionBehavior(DbContext dbContext, ILogger<TransactionBehavior<TRequest, TResponse>> logger)
{
_dbContext = dbContext;
_logger = logger;
}
public async Task<TResponse> Handle(
TRequest request,
RequestHandlerDelegate<TResponse> next,
CancellationToken cancellationToken)
{
// 只有写操作需要事务
if (typeof(TResponse) == typeof(Unit) || typeof(TResponse) == typeof(void))
{
var strategy = _dbContext.Database.CreateExecutionStrategy();
return await strategy.ExecuteAsync(async () =>
{
using var transaction = await _dbContext.Database.BeginTransactionAsync(cancellationToken);
try
{
var response = await next(cancellationToken);
await _dbContext.SaveChangesAsync(cancellationToken);
await transaction.CommitAsync(cancellationToken);
_logger.LogInformation("事务提交成功");
return response;
}
catch (Exception ex)
{
await transaction.RollbackAsync(cancellationToken);
_logger.LogError(ex, "事务回滚");
throw;
}
});
}
// 读操作直接执行
return await next(cancellationToken);
}
}注册
csharp
builder.Services.AddTransient(typeof(IPipelineBehavior<,>), typeof(TransactionBehavior<,>));💾 缓存行为(Caching Behavior)
实现
csharp
public class CachingBehavior<TRequest, TResponse> : IPipelineBehavior<TRequest, TResponse>
where TRequest : IRequest<TResponse>
{
private readonly IDistributedCache _cache;
private readonly ILogger<CachingBehavior<TRequest, TResponse>> _logger;
private readonly TimeSpan _defaultExpiration = TimeSpan.FromMinutes(5);
public CachingBehavior(IDistributedCache cache, ILogger<CachingBehavior<TRequest, TResponse>> logger)
{
_cache = cache;
_logger = logger;
}
public async Task<TResponse> Handle(
TRequest request,
RequestHandlerDelegate<TResponse> next,
CancellationToken cancellationToken)
{
// 只对查询启用缓存
if (!typeof(TRequest).Name.EndsWith("Query"))
{
return await next(cancellationToken);
}
var cacheKey = GenerateCacheKey(request);
// 尝试从缓存获取
var cachedResult = await _cache.GetStringAsync(cacheKey, cancellationToken);
if (!string.IsNullOrEmpty(cachedResult))
{
_logger.LogInformation("缓存命中: {CacheKey}", cacheKey);
return JsonSerializer.Deserialize<TResponse>(cachedResult);
}
_logger.LogInformation("缓存未命中,执行查询: {CacheKey}", cacheKey);
// 执行实际查询
var response = await next(cancellationToken);
// 存入缓存
var serializedResult = JsonSerializer.Serialize(response);
await _cache.SetStringAsync(cacheKey, serializedResult, new DistributedCacheEntryOptions
{
AbsoluteExpirationRelativeToNow = _defaultExpiration
}, cancellationToken);
return response;
}
private string GenerateCacheKey(TRequest request)
{
var requestType = typeof(TRequest).Name;
var requestHash = ComputeHash(JsonSerializer.Serialize(request));
return $"{requestType}:{requestHash}";
}
private string ComputeHash(string input)
{
using var sha256 = SHA256.Create();
var bytes = sha256.ComputeHash(Encoding.UTF8.GetBytes(input));
return Convert.ToBase64String(bytes);
}
}注册
csharp
builder.Services.AddDistributedMemoryCache(); // 或使用 Redis
builder.Services.AddTransient(typeof(IPipelineBehavior<,>), typeof(CachingBehavior<,>));📊 性能监控行为
实现
csharp
public class PerformanceBehavior<TRequest, TResponse> : IPipelineBehavior<TRequest, TResponse>
where TRequest : IRequest<TResponse>
{
private readonly Stopwatch _timer;
private readonly ILogger<PerformanceBehavior<TRequest, TResponse>> _logger;
private readonly int _thresholdMilliseconds = 500; // 慢查询阈值
public PerformanceBehavior(ILogger<PerformanceBehavior<TRequest, TResponse>> logger)
{
_timer = new Stopwatch();
_logger = logger;
}
public async Task<TResponse> Handle(
TRequest request,
RequestHandlerDelegate<TResponse> next,
CancellationToken cancellationToken)
{
_timer.Start();
var response = await next(cancellationToken);
_timer.Stop();
var elapsedMilliseconds = _timer.ElapsedMilliseconds;
var requestName = typeof(TRequest).Name;
if (elapsedMilliseconds > _thresholdMilliseconds)
{
_logger.LogWarning(
"⚠️ 慢查询检测: {RequestName} 耗时 {ElapsedMilliseconds}ms {@Request}",
requestName,
elapsedMilliseconds,
request);
}
else
{
_logger.LogDebug(
"✅ {RequestName} 耗时 {ElapsedMilliseconds}ms",
requestName,
elapsedMilliseconds);
}
return response;
}
}注册
csharp
builder.Services.AddTransient(typeof(IPipelineBehavior<,>), typeof(PerformanceBehavior<,>));🔧 行为的注册顺序与执行流程
注册顺序决定执行顺序
csharp
// Program.cs
builder.Services.AddMediatR(cfg =>
cfg.RegisterServicesFromAssembly(typeof(Program).Assembly));
// 按顺序注册行为(重要!)
builder.Services.AddTransient(typeof(IPipelineBehavior<,>), typeof(LoggingBehavior<,>)); // 1. 日志
builder.Services.AddTransient(typeof(IPipelineBehavior<,>), typeof(ValidationBehavior<,>)); // 2. 验证
builder.Services.AddTransient(typeof(IPipelineBehavior<,>), typeof(TransactionBehavior<,>)); // 3. 事务
builder.Services.AddTransient(typeof(IPipelineBehavior<,>), typeof(CachingBehavior<,>)); // 4. 缓存
builder.Services.AddTransient(typeof(IPipelineBehavior<,>), typeof(PerformanceBehavior<,>)); // 5. 性能监控执行流程图
mermaid
graph TB
A[Request] --> B[1. LoggingBehavior<br/>记录开始日志]
B --> C[2. ValidationBehavior<br/>验证参数]
C --> D{验证通过?}
D -->|否| E[抛出 ValidationException]
D -->|是| F[3. TransactionBehavior<br/>开启事务]
F --> G[4. CachingBehavior<br/>检查缓存]
G --> H{缓存命中?}
H -->|是| I[返回缓存数据]
H -->|否| J[5. PerformanceBehavior<br/>开始计时]
J --> K[Handler<br/>业务逻辑]
K --> L[5. PerformanceBehavior<br/>结束计时]
L --> M[4. CachingBehavior<br/>写入缓存]
M --> N[3. TransactionBehavior<br/>提交事务]
N --> B
B --> O[Response]
style E fill:#ff6b6b
style I fill:#51cf66
style K fill:#ffd43b洋葱模型
┌─────────────────────────────────────────┐
│ LoggingBehavior (最外层) │
│ ┌───────────────────────────────────┐ │
│ │ ValidationBehavior │ │
│ │ ┌─────────────────────────────┐ │ │
│ │ │ TransactionBehavior │ │ │
│ │ │ ┌───────────────────────┐ │ │ │
│ │ │ │ CachingBehavior │ │ │ │
│ │ │ │ ┌─────────────────┐ │ │ │ │
│ │ │ │ │ PerformanceBeh. │ │ │ │ │
│ │ │ │ │ ┌───────────┐ │ │ │ │ │
│ │ │ │ │ │ Handler │ │ │ │ │ │
│ │ │ │ │ └───────────┘ │ │ │ │ │
│ │ │ │ └─────────────────┘ │ │ │ │
│ │ │ └───────────────────────┘ │ │ │
│ │ └─────────────────────────────┘ │ │
│ └───────────────────────────────────┘ │
└─────────────────────────────────────────┘🎭 条件性行为
根据请求类型决定是否执行
csharp
public class ConditionalLoggingBehavior<TRequest, TResponse> : IPipelineBehavior<TRequest, TResponse>
where TRequest : IRequest<TResponse>
{
private readonly ILogger<ConditionalLoggingBehavior<TRequest, TResponse>> _logger;
public async Task<TResponse> Handle(
TRequest request,
RequestHandlerDelegate<TResponse> next,
CancellationToken cancellationToken)
{
// 只记录命令,不记录查询
if (typeof(TRequest).Name.EndsWith("Command"))
{
_logger.LogInformation("执行命令: {RequestName}", typeof(TRequest).Name);
}
return await next(cancellationToken);
}
}使用特性标记控制
csharp
// 自定义特性
[AttributeUsage(AttributeTargets.Class)]
public class SkipValidationAttribute : Attribute
{
}
// 标记跳过验证的请求
[SkipValidation]
public class BulkDeleteCommand : IRequest
{
public List<Guid> Ids { get; set; }
}
// 验证行为中检查特性
public class ValidationBehavior<TRequest, TResponse> : IPipelineBehavior<TRequest, TResponse>
{
public async Task<TResponse> Handle(
TRequest request,
RequestHandlerDelegate<TResponse> next,
CancellationToken cancellationToken)
{
// 检查是否有跳过验证的特性
var skipValidation = typeof(TRequest)
.GetCustomAttributes(typeof(SkipValidationAttribute), true)
.Any();
if (skipValidation)
{
return await next(cancellationToken);
}
// 正常验证逻辑
// ...
}
}🔄 多个行为的组合与调试
调试技巧:可视化执行顺序
csharp
public class DebugBehavior<TRequest, TResponse> : IPipelineBehavior<TRequest, TResponse>
{
private readonly ILogger<DebugBehavior<TRequest, TResponse>> _logger;
public async Task<TResponse> Handle(
TRequest request,
RequestHandlerDelegate<TResponse> next,
CancellationToken cancellationToken)
{
var behaviorName = GetType().Name;
var requestName = typeof(TRequest).Name;
_logger.LogInformation("➡️ 进入 {BehaviorName} 处理 {RequestName}", behaviorName, requestName);
try
{
var response = await next(cancellationToken);
_logger.LogInformation("⬅️ 离开 {BehaviorName} 处理 {RequestName}", behaviorName, requestName);
return response;
}
catch (Exception ex)
{
_logger.LogError(ex, "❌ {BehaviorName} 中发生异常", behaviorName);
throw;
}
}
}输出示例
➡️ 进入 LoggingBehavior 处理 CreateOrderCommand
➡️ 进入 ValidationBehavior 处理 CreateOrderCommand
➡️ 进入 TransactionBehavior 处理 CreateOrderCommand
➡️ 进入 CachingBehavior 处理 CreateOrderCommand
➡️ 进入 PerformanceBehavior 处理 CreateOrderCommand
[Handler 执行业务逻辑]
⬅️ 离开 PerformanceBehavior 处理 CreateOrderCommand
⬅️ 离开 CachingBehavior 处理 CreateOrderCommand
⬅️ 离开 TransactionBehavior 处理 CreateOrderCommand
⬅️ 离开 ValidationBehavior 处理 CreateOrderCommand
⬅️ 离开 LoggingBehavior 处理 CreateOrderCommand🎯 最佳实践
✅ 推荐做法
csharp
// 1. 保持 Behavior 单一职责
public class LoggingBehavior<TRequest, TResponse> : IPipelineBehavior<TRequest, TResponse>
{
// 只负责日志记录
}
// 2. 使用泛型约束限制适用范围
public class ValidationBehavior<TRequest, TResponse> : IPipelineBehavior<TRequest, TResponse>
where TRequest : IRequest<TResponse>
{
// 只处理 IRequest
}
// 3. 异步调用 next
public async Task<TResponse> Handle(TRequest request, RequestHandlerDelegate<TResponse> next, CancellationToken ct)
{
return await next(ct); // ✅ 正确
}
// 4. 正确处理异常
try
{
return await next(ct);
}
catch (Exception ex)
{
_logger.LogError(ex, "处理失败");
throw; // ✅ 重新抛出
}❌ 避免的做法
csharp
// 1. 不要在 Behavior 中执行耗时同步操作
public Task<TResponse> Handle(TRequest request, RequestHandlerDelegate<TResponse> next, CancellationToken ct)
{
Thread.Sleep(1000); // ❌ 阻塞线程
return next(ct);
}
// 2. 不要忘记调用 next
public Task<TResponse> Handle(TRequest request, RequestHandlerDelegate<TResponse> next, CancellationToken ct)
{
// 忘记调用 next() // ❌ Handler 永远不会执行
return Task.FromResult(default(TResponse));
}
// 3. 不要吞掉异常
try
{
return await next(ct);
}
catch (Exception ex)
{
_logger.LogError(ex, "处理失败");
return default; // ❌ 吞掉异常,会导致静默失败
}📊 性能考量
行为链的性能影响
| 行为数量 | 平均开销 | 说明 |
|---|---|---|
| 0 | 0ms | 无行为 |
| 1-2 | ~0.01ms | 可忽略 |
| 3-5 | ~0.05ms | 轻微影响 |
| 6-10 | ~0.1ms | 中等影响 |
| 10+ | ~0.2ms+ | 显著影响 |
结论:合理使用 5-10 个行为对性能影响很小,但应避免过多行为。
🎓 总结
核心价值
- 解耦横切关注点:日志、验证、事务等与业务逻辑分离
- 复用性强:一个行为自动应用于所有请求
- 灵活性高:可以按需组合、条件执行
- 易于测试:每个行为独立测试
常见应用场景
| 场景 | 推荐行为 |
|---|---|
| 日志记录 | LoggingBehavior |
| 参数验证 | ValidationBehavior |
| 事务管理 | TransactionBehavior |
| 数据缓存 | CachingBehavior |
| 性能监控 | PerformanceBehavior |
| 安全检查 | AuthorizationBehavior |
| 重试机制 | RetryBehavior |
| 限流控制 | RateLimitingBehavior |
🚀 下一步
- 通知机制 - 学习发布/订阅模式
- ASP.NET Core 集成 - 在 Web 应用中应用管道行为
- 实战案例 - 完整的 CQRS 系统中的管道行为应用
💡 提示:管道行为是 MediatR 的灵魂,掌握它就掌握了 MediatR 的精髓!