Skip to content

与 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 优势总结 ​

维度传统 CRUDCQRS
复杂度低中
读性能一般✅ 优秀
写性能一般✅ 优秀
可扩展性有限✅ 高
维护成本低中
适用场景简单应用复杂业务

🎓 总结 ​

CQRS + MediatR 的核心价值 ​

✅ 清晰的职责分离:命令和查询完全独立
✅ 灵活的技术选型:读写可使用不同技术栈
✅ 优异的性能:各自优化,互不影响
✅ 良好的扩展性:易于演变为微服务

何时使用 CQRS? ​

  • ✅ 读写比例悬殊(如 10:1)
  • ✅ 复杂的业务逻辑
  • ✅ 高性能要求
  • ✅ 团队协作的大型项目

💡 提示: CQRS 不是银弹,简单 CRUD 应用无需过度设计!

Released under the CC BY-SA 4.0 License.