Skip to content

管道行为(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);
}

核心概念 ​

参数类型说明
requestTRequest当前请求对象
nextRequestHandlerDelegate<TResponse>委托,指向下一个处理器(可能是另一个 Behavior 或最终的 Handler)
cancellationTokenCancellationToken取消令牌
返回值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]

关键点:

  1. 请求依次通过每个 Behavior(外层 → 内层)
  2. 到达最终的 Handler
  3. 响应依次返回(内层 → 外层)

📝 日志行为(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; // ❌ 吞掉异常,会导致静默失败
}

📊 性能考量 ​

行为链的性能影响 ​

行为数量平均开销说明
00ms无行为
1-2~0.01ms可忽略
3-5~0.05ms轻微影响
6-10~0.1ms中等影响
10+~0.2ms+显著影响

结论:合理使用 5-10 个行为对性能影响很小,但应避免过多行为。


🎓 总结 ​

核心价值 ​

  1. 解耦横切关注点:日志、验证、事务等与业务逻辑分离
  2. 复用性强:一个行为自动应用于所有请求
  3. 灵活性高:可以按需组合、条件执行
  4. 易于测试:每个行为独立测试

常见应用场景 ​

场景推荐行为
日志记录LoggingBehavior
参数验证ValidationBehavior
事务管理TransactionBehavior
数据缓存CachingBehavior
性能监控PerformanceBehavior
安全检查AuthorizationBehavior
重试机制RetryBehavior
限流控制RateLimitingBehavior

🚀 下一步 ​


💡 提示:管道行为是 MediatR 的灵魂,掌握它就掌握了 MediatR 的精髓!

Released under the CC BY-SA 4.0 License.