Skip to content

高级请求类型 ​

🚀 流式处理、无返回值请求和泛型约束


📖 概述 ​

除了标准的 IRequest<TResponse>,MediatR 还支持一些特殊的请求类型,用于满足不同的业务场景。


1️⃣ IStreamRequest<TResponse> - 流式处理 ​

使用场景 ​

  • 大数据集分页加载
  • 实时数据流
  • 文件逐行处理
  • SSE(Server-Sent Events)

定义流式请求 ​

csharp
public class StreamOrdersQuery : IStreamRequest<OrderDto>
{
    public Guid? CustomerId { get; set; }
    public int PageSize { get; set; } = 100;
}

实现流式处理器 ​

csharp
public class StreamOrdersHandler : IStreamRequestHandler<StreamOrdersQuery, OrderDto>
{
    private readonly IOrderRepository _orderRepo;

    public StreamOrdersHandler(IOrderRepository orderRepo)
    {
        _orderRepo = orderRepo;
    }

    public async IAsyncEnumerable<OrderDto> Handle(
        StreamOrdersQuery request, 
        [EnumeratorCancellation] CancellationToken cancellationToken)
    {
        var skip = 0;

        while (!cancellationToken.IsCancellationRequested)
        {
            // 分批获取数据
            var orders = await _orderRepo.GetPageAsync(
                skip, 
                request.PageSize, 
                request.CustomerId,
                cancellationToken);

            if (!orders.Any())
                break;

            // 逐个返回订单
            foreach (var order in orders)
            {
                yield return new OrderDto
                {
                    Id = order.Id,
                    ProductName = order.ProductName,
                    TotalAmount = order.TotalAmount,
                    CreatedAt = order.CreatedAt
                };
            }

            skip += request.PageSize;
        }
    }
}

消费流式数据 ​

控制台应用 ​

csharp
var query = new StreamOrdersQuery { PageSize = 50 };

await foreach (var order in mediator.CreateStream(query))
{
    Console.WriteLine($"订单: {order.Id}, 金额: {order.TotalAmount:C}");
}

ASP.NET Core API(SSE) ​

csharp
[HttpGet("stream")]
public async Task StreamOrders([FromQuery] Guid? customerId)
{
    Response.ContentType = "text/event-stream";
    
    var query = new StreamOrdersQuery 
    { 
        CustomerId = customerId,
        PageSize = 100 
    };

    await foreach (var order in _mediator.CreateStream(query, HttpContext.RequestAborted))
    {
        var json = JsonSerializer.Serialize(order);
        await Response.WriteAsync($"data: {json}\n\n");
        await Response.Body.FlushAsync();
    }
}

Blazor Server ​

csharp
@code {
    private List<OrderDto> _orders = new();

    protected override async Task OnInitializedAsync()
    {
        var query = new StreamOrdersQuery { PageSize = 50 };

        await foreach (var order in Mediator.CreateStream(query))
        {
            _orders.Add(order);
            StateHasChanged(); // 实时更新 UI
        }
    }
}

2️⃣ Unit / Void - 无返回值请求 ​

方式 1:使用 IRequest(推荐) ​

csharp
// 定义命令
public class DeleteOrderCommand : IRequest
{
    public Guid OrderId { get; set; }
}

// 实现处理器
public class DeleteOrderHandler : IRequestHandler<DeleteOrderCommand>
{
    private readonly IOrderRepository _orderRepo;

    public DeleteOrderHandler(IOrderRepository orderRepo)
    {
        _orderRepo = orderRepo;
    }

    public async Task Handle(DeleteOrderCommand request, CancellationToken cancellationToken)
    {
        await _orderRepo.DeleteAsync(request.OrderId, cancellationToken);
    }
}

// 发送命令
await _mediator.Send(new DeleteOrderCommand { OrderId = orderId });

方式 2:使用 Unit 类型 ​

Unit 是 MediatR 中表示"无值"的类型(类似于 void 但可用于泛型)。

csharp
public class LogActionCommand : IRequest<Unit>
{
    public string Action { get; set; }
}

public class LogActionHandler : IRequestHandler<LogActionCommand, Unit>
{
    public Task<Unit> Handle(LogActionCommand request, CancellationToken cancellationToken)
    {
        Console.WriteLine($"记录操作: {request.Action}");
        return Unit.Task; // 返回 Unit 实例
    }
}

// 发送
await _mediator.Send(new LogActionCommand { Action = "UserLogin" });

Unit vs IRequest 对比 ​

特性IRequestIRequest<Unit>
简洁性✅ 更简洁⚠️ 需要返回 Unit.Task
意图表达✅ 清晰表示无返回值⚠️ 不够直观
推荐程度✅ 推荐⚠️ 兼容旧代码

建议:优先使用 IRequest(无泛型参数)。


3️⃣ 请求约束与泛型处理 ​

泛型约束 ​

csharp
// 只处理特定类型的请求
public class AuditBehavior<TRequest, TResponse> : IPipelineBehavior<TRequest, TResponse>
    where TRequest : IAuditableRequest  // 自定义接口约束
{
    public async Task<TResponse> Handle(
        TRequest request, 
        RequestHandlerDelegate<TResponse> next, 
        CancellationToken cancellationToken)
    {
        // 只审计实现了 IAuditableRequest 的请求
        await _auditService.LogAsync(request);
        return await next(cancellationToken);
    }
}

// 标记接口
public interface IAuditableRequest
{
    Guid CorrelationId { get; set; }
}

// 实现接口
public class CreateOrderCommand : IRequest<OrderResult>, IAuditableRequest
{
    public Guid CorrelationId { get; set; } = Guid.NewGuid();
    public string ProductName { get; set; }
}

泛型处理器 ​

csharp
// 通用查询处理器(适用于所有 Query)
public class CachingQueryHandler<TQuery, TResponse> : IRequestHandler<TQuery, TResponse>
    where TQuery : ICacheableQuery<TResponse>
{
    private readonly IDistributedCache _cache;
    private readonly IRequestHandler<TQuery, TResponse> _innerHandler;

    public async Task<TResponse> Handle(TQuery request, CancellationToken cancellationToken)
    {
        var cacheKey = request.GetCacheKey();
        
        // 尝试从缓存获取
        var cached = await _cache.GetStringAsync(cacheKey, cancellationToken);
        if (!string.IsNullOrEmpty(cached))
            return JsonSerializer.Deserialize<TResponse>(cached);

        // 执行实际查询
        var result = await _innerHandler.Handle(request, cancellationToken);

        // 存入缓存
        await _cache.SetStringAsync(cacheKey, JsonSerializer.Serialize(result), 
            new DistributedCacheEntryOptions { AbsoluteExpirationRelativeToNow = TimeSpan.FromMinutes(5) },
            cancellationToken);

        return result;
    }
}

🎯 最佳实践 ​

✅ 推荐做法 ​

csharp
// 1. 流式处理用于大数据集
public IAsyncEnumerable<OrderDto> Handle(StreamOrdersQuery request, ...)
{
    // 分批加载,避免内存溢出
}

// 2. 无返回值使用 IRequest(非泛型)
public class DeleteCommand : IRequest { }

// 3. 使用接口约束实现条件行为
where TRequest : IAuditableRequest

❌ 避免的做法 ​

csharp
// 1. 不要在小数据集上使用流式处理
// 直接返回 List<T> 更简单

// 2. 不要过度使用 Unit
// 优先使用 IRequest(无泛型)

// 3. 不要在流式处理中一次性加载所有数据
public async IAsyncEnumerable<OrderDto> Handle(...)
{
    var allOrders = await _repo.GetAllAsync(); // ❌ 失去流式优势
    foreach (var order in allOrders)
        yield return order;
}

📊 性能对比 ​

场景传统方式流式处理性能提升
10万条记录内存占用 ~500MB内存占用 ~5MB✅ 99%
首字节时间等待全部加载立即返回✅ 10倍+
用户体验长时间等待渐进式加载✅ 显著提升

🎓 总结 ​

选择指南 ​

需求推荐类型
标准请求/响应IRequest<TResponse>
无返回值操作IRequest
大数据集/实时流IStreamRequest<TResponse>
条件性行为泛型约束 + 标记接口

💡 提示:流式处理是优化大数据场景的神器,合理使用可显著提升性能和用户体验!

Released under the CC BY-SA 4.0 License.