与 CQRS 架构集成
🏗️ 使用 MediatR 实现命令查询职责分离(CQRS)
📖 什么是 CQRS?
CQRS(Command Query Responsibility Segregation)是一种架构模式,核心思想是:
- 命令(Command): 修改状态的操作(写),不返回数据
- 查询(Query): 读取数据的操作(读),不修改状态
通过分离读写操作,可以:
- ✅ 独立优化读写性能
- ✅ 使用不同的数据模型
- ✅ 提高系统可扩展性
- ✅ 简化业务逻辑
🎯 命令与查询的区分
命名约定
csharp
// 命令(写操作)- 使用 Command 后缀
public class CreateOrderCommand : IRequest<OrderResult> { }
public class UpdateOrderCommand : IRequest { }
public class DeleteOrderCommand : IRequest { }
// 查询(读操作)- 使用 Query 后缀
public class GetOrderQuery : IRequest<OrderDto> { }
public class ListOrdersQuery : IRequest<List<OrderDto>> { }
public class OrderStatisticsQuery : IRequest<OrderStats> { }标记接口(可选)
csharp
// 标记接口,用于区分命令和查询
public interface ICommand : IRequest { }
public interface ICommand<TResponse> : IRequest<TResponse> { }
public interface IQuery<TResponse> : IRequest<TResponse> { }
// 使用标记接口
public class CreateOrderCommand : ICommand<OrderResult>
{
public string ProductName { get; set; }
}
public class GetOrderQuery : IQuery<OrderDto>
{
public Guid OrderId { get; set; }
}🔧 命令处理器与查询处理器分离
命令处理器(写)
csharp
public class CreateOrderHandler : IRequestHandler<CreateOrderCommand, OrderResult>
{
private readonly IOrderRepository _orderRepo;
private readonly IMediator _mediator;
public CreateOrderHandler(IOrderRepository orderRepo, IMediator mediator)
{
_orderRepo = orderRepo;
_mediator = mediator;
}
public async Task<OrderResult> Handle(CreateOrderCommand request, CancellationToken ct)
{
// 1. 业务验证
if (request.Quantity <= 0)
throw new ValidationException("数量必须大于0");
// 2. 创建订单(领域模型)
var order = new Order
{
Id = Guid.NewGuid(),
ProductName = request.ProductName,
Quantity = request.Quantity,
Price = request.Price,
CreatedAt = DateTime.UtcNow
};
// 3. 持久化
await _orderRepo.AddAsync(order, ct);
// 4. 发布领域事件
await _mediator.Publish(new OrderCreatedEvent
{
OrderId = order.Id,
TotalAmount = order.TotalAmount
}, ct);
return new OrderResult { OrderId = order.Id };
}
}查询处理器(读)
csharp
public class GetOrderHandler : IRequestHandler<GetOrderQuery, OrderDto>
{
private readonly IDbConnection _db; // 直接使用 Dapper 或其他轻量级 ORM
public GetOrderHandler(IDbConnection db)
{
_db = db;
}
public async Task<OrderDto> Handle(GetOrderQuery request, CancellationToken ct)
{
// 直接使用 SQL 查询,优化性能
const string sql = @"
SELECT o.Id, o.ProductName, o.Quantity, o.Price,
c.Name as CustomerName, c.Email as CustomerEmail
FROM Orders o
INNER JOIN Customers c ON o.CustomerId = c.Id
WHERE o.Id = @OrderId";
return await _db.QueryFirstOrDefaultAsync<OrderDto>(sql, new
{
request.OrderId
});
}
}关键区别:
- 命令侧:使用领域模型,保证数据一致性
- 查询侧:使用 DTO,直接 SQL 优化性能
💾 读写分离与最终一致性
分离数据库
csharp
// 写数据库(主库)
public class WriteDbContext : DbContext
{
public DbSet<Order> Orders { get; set; }
protected override void OnConfiguring(DbContextOptionsBuilder options)
{
options.UseSqlServer("Server=master-db;Database=Orders;...");
}
}
// 读数据库(从库,可以是 NoSQL)
public class ReadDbContext : DbContext
{
public DbSet<OrderReadModel> Orders { get; set; }
protected override void OnConfiguring(DbContextOptionsBuilder options)
{
// 可以使用 MongoDB、Elasticsearch 等
options.UseMongoDB("mongodb://read-db/orders");
}
}同步策略:通过领域事件
csharp
// 当订单创建后,异步同步到读数据库
public class SyncOrderToReadModelHandler : INotificationHandler<OrderCreatedEvent>
{
private readonly ReadDbContext _readDb;
public async Task Handle(OrderCreatedEvent notification, CancellationToken ct)
{
var readModel = new OrderReadModel
{
Id = notification.OrderId,
TotalAmount = notification.TotalAmount,
CreatedAt = DateTime.UtcNow
};
await _readDb.Orders.AddAsync(readModel, ct);
await _readDb.SaveChangesAsync(ct);
}
}最终一致性:写操作完成后,读模型会在短时间内(通常毫秒级)更新。
🌟 完整 CQRS 订单系统案例
项目结构
OrderSystem.CQRS/
├── Commands/
│ ├── CreateOrder/
│ │ ├── CreateOrderCommand.cs
│ │ ├── CreateOrderHandler.cs
│ │ └── CreateOrderValidator.cs
│ ├── UpdateOrder/
│ └── DeleteOrder/
├── Queries/
│ ├── GetOrder/
│ │ ├── GetOrderQuery.cs
│ │ └── GetOrderHandler.cs
│ ├── ListOrders/
│ └── OrderStatistics/
├── Events/
│ ├── OrderCreatedEvent.cs
│ ├── OrderUpdatedEvent.cs
│ └── Handlers/
├── Models/
│ ├── Order.cs (领域模型)
│ └── OrderDto.cs (DTO)
└── Controllers/
└── OrdersController.cs命令示例
csharp
// Commands/CreateOrder/CreateOrderCommand.cs
public class CreateOrderCommand : ICommand<OrderResult>
{
public Guid CustomerId { get; set; }
public List<OrderItemCommand> Items { get; set; }
public string ShippingAddress { get; set; }
}
public class OrderItemCommand
{
public Guid ProductId { get; set; }
public int Quantity { get; set; }
public decimal Price { get; set; }
}
// Commands/CreateOrder/CreateOrderHandler.cs
public class CreateOrderHandler : IRequestHandler<CreateOrderCommand, OrderResult>
{
private readonly WriteDbContext _writeDb;
private readonly IMediator _mediator;
public async Task<OrderResult> Handle(CreateOrderCommand request, CancellationToken ct)
{
// 创建聚合根
var order = Order.Create(
request.CustomerId,
request.Items.Select(i => new OrderItem(i.ProductId, i.Quantity, i.Price)),
request.ShippingAddress
);
// 保存
await _writeDb.Orders.AddAsync(order, ct);
await _writeDb.SaveChangesAsync(ct);
// 发布事件
await _mediator.Publish(new OrderCreatedEvent
{
OrderId = order.Id,
CustomerId = order.CustomerId,
TotalAmount = order.TotalAmount
}, ct);
return new OrderResult { OrderId = order.Id };
}
}查询示例
csharp
// Queries/ListOrders/ListOrdersQuery.cs
public class ListOrdersQuery : IQuery<PagedResult<OrderSummary>>
{
public Guid? CustomerId { get; set; }
public int PageNumber { get; set; } = 1;
public int PageSize { get; set; } = 20;
}
// Queries/ListOrders/ListOrdersHandler.cs
public class ListOrdersHandler : IRequestHandler<ListOrdersQuery, PagedResult<OrderSummary>>
{
private readonly ReadDbContext _readDb;
public async Task<PagedResult<OrderSummary>> Handle(ListOrdersQuery request, CancellationToken ct)
{
var query = _readDb.Orders.AsQueryable();
// 过滤
if (request.CustomerId.HasValue)
query = query.Where(o => o.CustomerId == request.CustomerId.Value);
// 总数
var total = await query.CountAsync(ct);
// 分页
var orders = await query
.OrderByDescending(o => o.CreatedAt)
.Skip((request.PageNumber - 1) * request.PageSize)
.Take(request.PageSize)
.Select(o => new OrderSummary
{
Id = o.Id,
TotalAmount = o.TotalAmount,
Status = o.Status,
CreatedAt = o.CreatedAt
})
.ToListAsync(ct);
return new PagedResult<OrderSummary>
{
Items = orders,
TotalCount = total,
PageNumber = request.PageNumber,
PageSize = request.PageSize
};
}
}控制器
csharp
[ApiController]
[Route("api/[controller]")]
public class OrdersController : ControllerBase
{
private readonly IMediator _mediator;
public OrdersController(IMediator mediator)
{
_mediator = mediator;
}
// 写操作 - 命令
[HttpPost]
public async Task<ActionResult<OrderResult>> CreateOrder([FromBody] CreateOrderCommand command)
{
var result = await _mediator.Send(command);
return CreatedAtAction(nameof(GetOrder), new { id = result.OrderId }, result);
}
// 读操作 - 查询
[HttpGet("{id}")]
public async Task<ActionResult<OrderDto>> GetOrder(Guid id)
{
var query = new GetOrderQuery { OrderId = id };
var order = await _mediator.Send(query);
if (order == null)
return NotFound();
return Ok(order);
}
// 列表查询
[HttpGet]
public async Task<ActionResult<PagedResult<OrderSummary>>> ListOrders(
[FromQuery] Guid? customerId,
[FromQuery] int page = 1,
[FromQuery] int pageSize = 20)
{
var query = new ListOrdersQuery
{
CustomerId = customerId,
PageNumber = page,
PageSize = pageSize
};
var result = await _mediator.Send(query);
return Ok(result);
}
}🎯 最佳实践
✅ 推荐做法
csharp
// 1. 命令和查询使用不同的数据模型
// 命令:领域模型(Order)
// 查询:DTO(OrderDto)
// 2. 查询侧使用轻量级 ORM(Dapper)
public class GetOrderHandler : IRequestHandler<GetOrderQuery, OrderDto>
{
private readonly IDbConnection _db; // Dapper
public async Task<OrderDto> Handle(GetOrderQuery request, CancellationToken ct)
{
return await _db.QueryFirstOrDefaultAsync<OrderDto>(
"SELECT * FROM Orders WHERE Id = @Id",
new { request.OrderId });
}
}
// 3. 使用管道行为统一处理
builder.Services.AddTransient(typeof(IPipelineBehavior<,>), typeof(ValidationBehavior<,>)); // 验证命令
builder.Services.AddTransient(typeof(IPipelineBehavior<,>), typeof(TransactionBehavior<,>)); // 事务❌ 避免的做法
csharp
// 1. 不要在查询中使用领域模型
public class GetOrderHandler : IRequestHandler<GetOrderQuery, Order> // ❌ 应该用 DTO
{
}
// 2. 不要在命令中返回大量数据
public class CreateOrderCommand : IRequest<OrderDetails> // ❌ 只返回 ID 或简单结果
{
}
// 3. 不要混用读写数据库
public class GetOrderHandler
{
private readonly WriteDbContext _db; // ❌ 查询应该用读数据库
}📊 CQRS 优势总结
| 维度 | 传统 CRUD | CQRS |
|---|---|---|
| 复杂度 | 低 | 中 |
| 读性能 | 一般 | ✅ 优秀 |
| 写性能 | 一般 | ✅ 优秀 |
| 可扩展性 | 有限 | ✅ 高 |
| 维护成本 | 低 | 中 |
| 适用场景 | 简单应用 | 复杂业务 |
🎓 总结
CQRS + MediatR 的核心价值
✅ 清晰的职责分离:命令和查询完全独立
✅ 灵活的技术选型:读写可使用不同技术栈
✅ 优异的性能:各自优化,互不影响
✅ 良好的扩展性:易于演变为微服务
何时使用 CQRS?
- ✅ 读写比例悬殊(如 10:1)
- ✅ 复杂的业务逻辑
- ✅ 高性能要求
- ✅ 团队协作的大型项目
💡 提示: CQRS 不是银弹,简单 CRUD 应用无需过度设计!