高级请求类型
🚀 流式处理、无返回值请求和泛型约束
📖 概述
除了标准的 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 对比
| 特性 | IRequest | IRequest<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> |
| 条件性行为 | 泛型约束 + 标记接口 |
💡 提示:流式处理是优化大数据场景的神器,合理使用可显著提升性能和用户体验!