一、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需要权衡一致性和性能,采用事件驱动的方式实现数据同步。