📖 数据密集型设计

CQRS命令查询职责分离

深入探讨CQRS架构模式与读写分离策略

一、CQRS概述

CQRS(Command Query Responsibility Segregation)是一种架构模式,将系统操作分为命令和查询两部分。命令负责数据写入和状态变更,查询负责数据读取。

二、CQRS架构

2.1 CQRS架构图

graph TD A[客户端] --> B[命令接口] A --> C[查询接口] B --> D[命令处理程序] D --> E[领域模型] E --> F[写数据库] C --> G[查询处理程序] G --> H[读数据库] F --> I[数据同步] I --> H E --> J[事件发布] J --> K[事件存储] K --> L[事件处理器] L --> H

2.2 CQRS核心概念

概念 描述 特点
命令 改变状态的操作 不返回数据、可能失败
查询 读取数据的操作 返回数据、不改变状态
写模型 处理命令的模型 事务性、一致性
读模型 处理查询的模型 优化查询、非规范化
事件 状态变更的记录 不可变、可追溯

三、命令处理

3.1 命令定义

public interface ICommand
{
    Guid Id { get; }
    DateTime Timestamp { get; }
}

public interface ICommandHandler<in TCommand> where TCommand : ICommand
{
    Task<CommandResult> HandleAsync(TCommand command);
}

public class CommandResult
{
    public bool Success { get; set; }
    public string Message { get; set; }
    public Guid? AggregateId { get; set; }
    public List<string> Errors { get; set; } = new List<string>();
}

public class CreateOrderCommand : ICommand
{
    public Guid Id { get; } = Guid.NewGuid();
    public DateTime Timestamp { get; } = DateTime.UtcNow;
    public Guid UserId { get; set; }
    public List<OrderItemDto> Items { get; set; }
    public string ShippingAddress { get; set; }
}

public class UpdateOrderCommand : ICommand
{
    public Guid Id { get; } = Guid.NewGuid();
    public DateTime Timestamp { get; } = DateTime.UtcNow;
    public Guid OrderId { get; set; }
    public OrderStatus NewStatus { get; set; }
}

3.2 命令处理器

public class CreateOrderCommandHandler : ICommandHandler<CreateOrderCommand>
{
    private readonly IOrderRepository _orderRepository;
    private readonly IProductRepository _productRepository;
    private readonly IEventBus _eventBus;
    
    public async Task<CommandResult> HandleAsync(CreateOrderCommand command)
    {
        var products = await _productRepository.GetByIdsAsync(command.Items.Select(i => i.ProductId));
        
        var order = Order.Create(command.UserId, command.Items, products, command.ShippingAddress);
        
        foreach (var validationError in order.Validate())
        {
            return CommandResult.Failure(validationError);
        }
        
        await _orderRepository.SaveAsync(order);
        
        await _eventBus.PublishAsync(new OrderCreatedEvent(order.Id, command.UserId));
        
        return CommandResult.Success(order.Id);
    }
}

public class OrderCommandHandler : ICommandHandler<UpdateOrderCommand>
{
    private readonly IOrderRepository _orderRepository;
    private readonly IEventBus _eventBus;
    
    public async Task<CommandResult> HandleAsync(UpdateOrderCommand command)
    {
        var order = await _orderRepository.GetByIdAsync(command.OrderId);
        
        if (order == null)
            return CommandResult.Failure("订单不存在");
        
        order.UpdateStatus(command.NewStatus);
        
        await _orderRepository.SaveAsync(order);
        
        await _eventBus.PublishAsync(new OrderStatusChangedEvent(order.Id, command.NewStatus));
        
        return CommandResult.Success(order.Id);
    }
}

3.3 命令总线

public class CommandBus
{
    private readonly IServiceProvider _serviceProvider;
    
    public async Task<CommandResult> SendAsync<TCommand>(TCommand command) where TCommand : ICommand
    {
        var handler = _serviceProvider.GetService<ICommandHandler<TCommand>>();
        
        if (handler == null)
            throw new InvalidOperationException($"No handler found for command {typeof(TCommand).Name}");
        
        return await handler.HandleAsync(command);
    }
    
    public async Task<CommandResult> SendAsync(ICommand command)
    {
        var handlerType = typeof(ICommandHandler<>).MakeGenericType(command.GetType());
        var handler = _serviceProvider.GetService(handlerType);
        
        if (handler == null)
            throw new InvalidOperationException($"No handler found for command {command.GetType().Name}");
        
        var method = handlerType.GetMethod("HandleAsync");
        var result = await (Task<CommandResult>)method.Invoke(handler, new object[] { command });
        
        return result;
    }
}

四、查询处理

4.1 查询定义

public interface IQuery<out TResult>
{
}

public interface IQueryHandler<in TQuery, out TResult> where TQuery : IQuery<TResult>
{
    Task<TResult> HandleAsync(TQuery query);
}

public class GetOrderQuery : IQuery<OrderDto>
{
    public Guid OrderId { get; set; }
}

public class GetOrdersByUserQuery : IQuery<List<OrderSummaryDto>>
{
    public Guid UserId { get; set; }
    public int Page { get; set; }
    public int PageSize { get; set; }
}

public class SearchOrdersQuery : IQuery<List<OrderSummaryDto>>
{
    public string Keyword { get; set; }
    public OrderStatus? Status { get; set; }
    public DateTime? StartDate { get; set; }
    public DateTime? EndDate { get; set; }
}

4.2 查询处理器

public class GetOrderQueryHandler : IQueryHandler<GetOrderQuery, OrderDto>
{
    private readonly IOrderReadRepository _orderReadRepository;
    
    public async Task<OrderDto> HandleAsync(GetOrderQuery query)
    {
        return await _orderReadRepository.GetByIdAsync(query.OrderId);
    }
}

public class GetOrdersByUserQueryHandler : IQueryHandler<GetOrdersByUserQuery, List<OrderSummaryDto>>
{
    private readonly IOrderReadRepository _orderReadRepository;
    
    public async Task<List<OrderSummaryDto>> HandleAsync(GetOrdersByUserQuery query)
    {
        return await _orderReadRepository.GetByUserIdAsync(
            query.UserId, 
            query.Page, 
            query.PageSize
        );
    }
}

public class SearchOrdersQueryHandler : IQueryHandler<SearchOrdersQuery, List<OrderSummaryDto>>
{
    private readonly IOrderReadRepository _orderReadRepository;
    
    public async Task<List<OrderSummaryDto>> HandleAsync(SearchOrdersQuery query)
    {
        return await _orderReadRepository.SearchAsync(
            query.Keyword, 
            query.Status, 
            query.StartDate, 
            query.EndDate
        );
    }
}

4.3 查询总线

public class QueryBus
{
    private readonly IServiceProvider _serviceProvider;
    
    public async Task<TResult> QueryAsync<TResult>(IQuery<TResult> query)
    {
        var handlerType = typeof(IQueryHandler<,>).MakeGenericType(query.GetType(), typeof(TResult));
        var handler = _serviceProvider.GetService(handlerType);
        
        if (handler == null)
            throw new InvalidOperationException($"No handler found for query {query.GetType().Name}");
        
        var method = handlerType.GetMethod("HandleAsync");
        var result = await (Task<TResult>)method.Invoke(handler, new object[] { query });
        
        return result;
    }
}

五、读模型与写模型

5.1 读写模型分离

graph TD A[写模型] --> B[领域实体] B --> C[业务逻辑] C --> D[事务处理] D --> E[写数据库] F[读模型] --> G[查询优化] G --> H[非规范化] H --> I[缓存] I --> J[读数据库] E --> K[数据同步] K --> J

5.2 写模型设计

public class Order : AggregateRoot
{
    private Guid _userId;
    private List<OrderItem> _items = new List<OrderItem>();
    private OrderStatus _status;
    private string _shippingAddress;
    private decimal _totalAmount;
    
    private Order() { }
    
    public static Order Create(Guid userId, List<OrderItemDto> items, 
        List<Product> products, string shippingAddress)
    {
        var order = new Order
        {
            _userId = userId,
            _shippingAddress = shippingAddress,
            _status = OrderStatus.Pending
        };
        
        foreach (var item in items)
        {
            var product = products.First(p => p.Id == item.ProductId);
            order._items.Add(new OrderItem(product.Id, product.Name, item.Quantity, product.Price));
        }
        
        order.CalculateTotalAmount();
        
        order.AddDomainEvent(new OrderCreatedEvent(order.Id, userId));
        
        return order;
    }
    
    public void UpdateStatus(OrderStatus newStatus)
    {
        if (!CanTransitionTo(newStatus))
            throw new InvalidOperationException($"Cannot transition from {_status} to {newStatus}");
        
        _status = newStatus;
        
        AddDomainEvent(new OrderStatusChangedEvent(Id, newStatus));
    }
    
    private void CalculateTotalAmount()
    {
        _totalAmount = _items.Sum(i => i.TotalPrice);
    }
    
    public List<string> Validate()
    {
        var errors = new List<string>();
        
        if (_items.Count == 0)
            errors.Add("订单不能为空");
        
        if (string.IsNullOrEmpty(_shippingAddress))
            errors.Add("收货地址不能为空");
        
        return errors;
    }
}

5.3 读模型设计

public class OrderDto
{
    public Guid Id { get; set; }
    public Guid UserId { get; set; }
    public string UserName { get; set; }
    public List<OrderItemDto> Items { get; set; }
    public OrderStatus Status { get; set; }
    public string StatusName { get; set; }
    public string ShippingAddress { get; set; }
    public decimal TotalAmount { get; set; }
    public DateTime CreatedAt { get; set; }
}

public class OrderSummaryDto
{
    public Guid Id { get; set; }
    public string OrderNumber { get; set; }
    public OrderStatus Status { get; set; }
    public string StatusName { get; set; }
    public decimal TotalAmount { get; set; }
    public int ItemCount { get; set; }
    public DateTime CreatedAt { get; set; }
}

public class OrderReadModel
{
    public Guid Id { get; set; }
    public Guid UserId { get; set; }
    public string UserName { get; set; }
    public int Status { get; set; }
    public string StatusName { get; set; }
    public string ShippingAddress { get; set; }
    public decimal TotalAmount { get; set; }
    public int ItemCount { get; set; }
    public DateTime CreatedAt { get; set; }
    public DateTime UpdatedAt { get; set; }
}

六、数据同步

6.1 数据同步策略

策略 描述 延迟 复杂度
同步复制 命令处理后立即更新读模型
事件驱动 通过事件更新读模型
CDC 基于数据库变更捕获
定时同步 定时任务同步数据

6.2 事件驱动同步

public class OrderCreatedEventHandler : IEventHandler<OrderCreatedEvent>
{
    private readonly IOrderReadRepository _orderReadRepository;
    
    public async Task HandleAsync(OrderCreatedEvent @event)
    {
        var order = await _orderRepository.GetByIdAsync(@event.OrderId);
        var user = await _userRepository.GetByIdAsync(@event.UserId);
        
        var readModel = new OrderReadModel
        {
            Id = order.Id,
            UserId = order.UserId,
            UserName = user.Name,
            Status = (int)order.Status,
            StatusName = order.Status.ToString(),
            ShippingAddress = order.ShippingAddress,
            TotalAmount = order.TotalAmount,
            ItemCount = order.Items.Count,
            CreatedAt = order.CreatedAt,
            UpdatedAt = order.UpdatedAt
        };
        
        await _orderReadRepository.SaveAsync(readModel);
    }
}

public class OrderStatusChangedEventHandler : IEventHandler<OrderStatusChangedEvent>
{
    private readonly IOrderReadRepository _orderReadRepository;
    
    public async Task HandleAsync(OrderStatusChangedEvent @event)
    {
        var readModel = await _orderReadRepository.GetByIdAsync(@event.OrderId);
        
        if (readModel != null)
        {
            readModel.Status = (int)@event.NewStatus;
            readModel.StatusName = @event.NewStatus.ToString();
            readModel.UpdatedAt = DateTime.UtcNow;
            
            await _orderReadRepository.UpdateAsync(readModel);
        }
    }
}

6.3 CDC同步

public class CdcSyncService
{
    public async Task SyncChangesAsync()
    {
        var changes = await _cdcClient.GetChangesAsync(_lastSyncTimestamp);
        
        foreach (var change in changes)
        {
            await ApplyChange(change);
        }
        
        _lastSyncTimestamp = DateTime.UtcNow;
    }
    
    private async Task ApplyChange(CdcChange change)
    {
        switch (change.OperationType)
        {
            case OperationType.Insert:
                await HandleInsert(change);
                break;
            case OperationType.Update:
                await HandleUpdate(change);
                break;
            case OperationType.Delete:
                await HandleDelete(change);
                break;
        }
    }
    
    private async Task HandleInsert(CdcChange change)
    {
        var readModel = MapToReadModel(change.Data);
        await _readRepository.InsertAsync(readModel);
    }
    
    private async Task HandleUpdate(CdcChange change)
    {
        var readModel = await _readRepository.GetByIdAsync(change.Key);
        UpdateReadModel(readModel, change.Data);
        await _readRepository.UpdateAsync(readModel);
    }
}

七、CQRS优势与挑战

7.1 CQRS优势

优势 描述
读写分离 读和写可以独立扩展
查询优化 读模型可以非规范化优化查询
关注点分离 命令和查询逻辑分离
可扩展性 读写可以使用不同的技术栈
事务边界清晰 命令处理有明确的事务边界

7.2 CQRS挑战

挑战 描述 解决方案
数据一致性 读写模型数据可能不一致 最终一致性、事件驱动
系统复杂度 增加了系统复杂度 合理设计、使用成熟框架
延迟 读模型更新有延迟 优化同步机制
维护成本 需要维护两套模型 自动化工具

八、CQRS适用场景

8.1 适用场景

graph TD A[CQRS适用场景] --> B[读写比例差异大] A --> C[复杂业务逻辑] A --> D[需要多维度查询] A --> E[需要事件溯源] A --> F[需要独立扩展] B --> B1[读多写少] C --> C1[复杂规则] D --> D1[报表、分析] E --> E1[审计、追溯] F --> F1[弹性伸缩]

8.2 不适用场景

  • 简单CRUD应用
  • 需要强一致性的场景
  • 团队规模小、项目时间紧
  • 数据量小、并发低

九、CQRS最佳实践

9.1 从简单开始

先实现简单的CQRS,再逐步引入事件溯源等高级特性。

9.2 使用成熟框架

使用MediatR等成熟框架实现命令和查询总线。

9.3 合理设计读模型

读模型应该根据查询需求进行非规范化设计。

9.4 处理一致性

采用最终一致性,使用事件驱动同步。

9.5 监控和日志

监控命令和查询的执行情况,记录详细日志。

十、总结

CQRS是一种强大的架构模式,通过分离命令和查询职责,能够提高系统的可扩展性和查询性能。CQRS适合读多写少、需要复杂业务逻辑和多维度查询的场景。实现CQRS需要权衡一致性和性能,采用事件驱动的方式实现数据同步。