在CQRS架构下,读写分离带来了性能提升,但也让写模型和读模型之间的数据一致性成了开发者最头疼的问题。很多朋友刚接触CQRS时,都会问:“我更新了数据,查询那边怎么还是旧的?”其实这就像你在厨房做好了一道菜(写操作),但餐厅的菜单(读模型)还没来得及更新,客人点的还是旧菜。别急,下面咱们就聊聊怎么解决这个“上菜延迟”的问题。
一、CQRS到底是怎么回事?
CQRS全称是命令查询职责分离。说白了,就是把“写数据”和“读数据”拆成两个独立的模型。写模型负责处理增删改(命令),读模型负责展示查询(查询)。这样做的好处是:写操作可以专心保证数据完整性和业务规则,读操作可以针对查询需求做优化,比如创建冗余的视图、使用缓存等。
但问题也来了:两个模型的数据源可能不一样。写模型更新了数据库,读模型可能从另一个缓存或专门查询库中读取,数据更新需要同步。同步过程中如果没有处理好,就会出现短期不一致——用户刚提交了订单,刷新页面却看不到新订单。
二、数据不一致是怎么发生的?
举个例子,你做了一个电商系统。用户下单时,写模型把订单保存到主库(写库),同时需要把一个“订单新增”事件推送到消息队列。读模型监听这个事件,然后更新自己专门用于展示的读库(比如Elasticsearch或者Redis)。如果消息队列有延迟,或者事件处理失败,用户查询订单列表时就会查不到刚刚下的单。
这种不一致在很多场景下是可以接受的,比如订单列表等个几秒再刷新就对了。但有些场景不行,比如支付扣款、库存扣减,用户需要立刻看到结果。所以我们需要根据业务容忍度选择不同的解决方案。
三、解决一致性的几种主流策略
3.1 最终一致性:接受短暂的延迟
这是最常用的策略。读模型不要求每时每刻都和写模型完全一致,但保证在一段时间后(比如几秒或几分钟)最终一致。具体做法是:写模型在更新主库后,发送领域事件到消息队列(比如RabbitMQ、Kafka),读模型订阅事件并更新自己的存储。
优点:实现简单,写模型性能不受影响,读模型可以高度优化。
缺点:用户可能在短时间内看到不一致数据,不适合强一致性要求的业务。
3.2 同步更新:宁可慢,不能错
如果业务要求写操作后立刻查询必须看到最新数据,可以让写模型在更新主库后,直接同步更新读模型。比如使用同一个数据库事务,同时写入写表和读表;或者使用内存缓存,写操作后同步清除或更新缓存。
优点:强一致性,用户感知无延迟。
缺点:写操作性能下降,读模型失去独立扩展能力,如果读模型出故障会阻塞写操作。
3.3 事件溯源:记录历史,回放真相
事件溯源是CQRS的绝配。写模型不保存当前状态,只记录发生了哪些事件(比如“订单已创建”“商品已加购物车”)。读模型从事件流中投影出需要的状态。由于事件是不可变的,只要事件顺序正确,读模型可以随时重放事件得到最新状态。
优点:天然解决一致性,因为读模型总是基于事件流重建,不会丢失任何变更;方便审计和回溯。
缺点:学习曲线陡峭,存储量大,查询复杂(需要投影)。
3.4 补偿机制:处理同步失败
即使采用了最终一致性,也可能出现事件丢失或处理失败的情况。这时需要一个补偿机制:比如定时扫描不一致记录,重试失败的事件;或者提供人工修复接口。也可以使用“幂等性”设计,让重复消费事件不会产生副作用。
四、实战:用C#实现一个简单的CQRS数据同步
我们以 .NET Core 为例,使用 MediatR 处理命令和事件,用 EF Core 操作数据库。假设有一个订单系统,用户下单后,写模型保存订单,然后发布事件,读模型监听事件并更新一个用于展示的读表(OrderReadModel)。
4.1 项目结构
Solution/
├── WriteModel/
│ ├── Commands/
│ ├── Events/
│ └── Handlers/
├── ReadModel/
│ ├── EventHandlers/
│ └── Queries/
└── Shared/
└── Events/
4.2 定义领域事件
// Shared/Events/OrderCreatedEvent.cs
public class OrderCreatedEvent : INotification
{
public Guid OrderId { get; set; }
public string CustomerName { get; set; }
public Decimal TotalAmount { get; set; }
public DateTime CreatedAt { get; set; }
}
4.3 写模型:处理命令并发布事件
// WriteModel/Commands/CreateOrderCommand.cs
public class CreateOrderCommand : IRequest<Guid>
{
public string CustomerName { get; set; }
public List<OrderItem> Items { get; set; }
}
// WriteModel/Commands/CreateOrderCommandHandler.cs
public class CreateOrderCommandHandler : IRequestHandler<CreateOrderCommand, Guid>
{
private readonly AppDbContext _writeDb;
private readonly IMediator _mediator;
public CreateOrderCommandHandler(AppDbContext writeDb, IMediator mediator)
{
_writeDb = writeDb;
_mediator = mediator;
}
public async Task<Guid> Handle(CreateOrderCommand request, CancellationToken cancellationToken)
{
// 1. 创建订单实体(写模型)
var order = new Order
{
Id = Guid.NewGuid(),
CustomerName = request.CustomerName,
TotalAmount = request.Items.Sum(i => i.Price * i.Quantity),
CreatedAt = DateTime.UtcNow
};
_writeDb.Orders.Add(order);
await _writeDb.SaveChangesAsync(cancellationToken); // 先保存到写库
// 2. 发布领域事件
await _mediator.Publish(new OrderCreatedEvent
{
OrderId = order.Id,
CustomerName = order.CustomerName,
TotalAmount = order.TotalAmount,
CreatedAt = order.CreatedAt
}, cancellationToken);
return order.Id;
}
}
4.4 读模型:订阅事件并更新读库
// ReadModel/EventHandlers/OrderCreatedEventHandler.cs
public class OrderCreatedEventHandler : INotificationHandler<OrderCreatedEvent>
{
private readonly ReadDbContext _readDb;
public OrderCreatedEventHandler(ReadDbContext readDb)
{
_readDb = readDb;
}
public async Task Handle(OrderCreatedEvent notification, CancellationToken cancellationToken)
{
// 将事件转换为读模型需要的格式
var orderReadModel = new OrderReadModel
{
Id = notification.OrderId,
CustomerName = notification.CustomerName,
TotalAmount = notification.TotalAmount,
CreatedAt = notification.CreatedAt,
Status = "已创建"
};
// 插入读库(可以是独立的数据库或表)
_readDb.OrderReadModels.Add(orderReadModel);
await _readDb.SaveChangesAsync(cancellationToken);
}
}
4.5 查库时的查询操作
// ReadModel/Queries/GetOrdersQuery.cs
public class GetOrdersQuery : IRequest<List<OrderReadModel>> { }
// ReadModel/Queries/GetOrdersQueryHandler.cs
public class GetOrdersQueryHandler : IRequestHandler<GetOrdersQuery, List<OrderReadModel>>
{
private readonly ReadDbContext _readDb;
public GetOrdersQueryHandler(ReadDbContext readDb)
{
_readDb = readDb;
}
public async Task<List<OrderReadModel>> Handle(GetOrdersQuery request, CancellationToken cancellationToken)
{
// 直接查询读库,不经过写模型
return await _readDb.OrderReadModels.ToListAsync(cancellationToken);
}
}
4.6 如何处理并发和失败?
在实际生产环境中,我们需要确保事件处理的可靠性。常见做法是:
- 使用消息队列(如RabbitMQ)代替进程内中介者,这样即使读模型服务重启,消息也不会丢失。
- 在事件处理中加入重试逻辑(比如Polly库)。
- 事件处理器实现幂等性:例如在OrderReadModel表中用OrderId作为唯一键,重复插入时忽略或更新。
五、应用场景分析
CQRS读写模型分离带来的数据一致性问题,并不是所有项目都适合用最终一致性。我们来梳理一下不同场景的选择:
适合最终一致性的场景:
- 内容管理系统(CMS):文章发布后,读者看到最新的内容可以接受几秒延迟。
- 社交动态:点赞、评论数更新不需要即时同步。
- 报表系统:数据可以延迟几小时。
需要强一致性的场景:
- 金融交易:账户余额扣减必须立即看到结果。
- 库存扣减:超卖会直接导致损失。
- 订单支付状态:用户付款后需要立即知道是否成功。
对于强一致性需求,可以使用同步更新或者事件溯源结合同一个事务。但要注意,同步更新会牺牲写性能,并且耦合了读模型,失去了CQRS的独立扩展优势。事件溯源虽然强一致,但实现复杂。
六、技术优缺点总结
最终一致性方案:
- 优点:写模型性能高,读模型可独立扩展、灵活优化;系统整体可用性强。
- 缺点:存在短暂不一致窗口;需要处理事件丢失、重复等异常;调试困难。
同步更新方案:
- 优点:读写强一致,代码逻辑直观。
- 缺点:写操作变慢,读模型故障会影响写操作;不能发挥CQRS的分离优势。
事件溯源方案:
- 优点:最强的可追溯性和一致性;支持时间旅行和审计。
- 缺点:学习成本高;查询需要投影,可能需要额外的专门读库;存储量大。
七、注意事项
- 不要为了用而用:如果你的系统读写量不大,或者读写模型一样,CQRS带来的复杂度可能大于收益。
- 事件总线的高可用:事件传递是整个系统的命脉,务必做好消息队列的集群、持久化和重试。
- 监控与补偿:建立一致性检查机制,比如定期对比写库和读库的数据差异,自动修复遗漏的事件。
- 领域事件设计:事件应该包含足够的信息,让读模型不需要再去查写库。例如订单创建事件应该包含所有需要展示的字段。
- 事务边界:在写模型中,发布事件和保存写库最好不要放在同一个本地事务里,因为如果事件发送失败,本地事务回滚会导致用户看到“操作成功”但实际没数据。更好的做法是先保存写库,然后发事件;如果事件发送失败,通过补偿任务重新发送。
八、文章总结
CQRS架构下的数据一致性问题,本质上是在“一致性”和“性能”之间做权衡。没有银弹,只能根据业务选择合适的策略。最终一致性适合大多数互联网应用,通过事件驱动的异步同步,可以同时获得高吞吐和较好的用户体验。对于核心的、强一致性的场景,可以考虑同步更新或事件溯源,但需要接受额外的复杂度。在实战中,一定要做好事件可靠传递、幂等处理和补偿机制。不要害怕不一致,但要对不一致窗口有清晰的认知,并告知业务方。记住,技术是为业务服务的,一致性方案的选择最终要回到“用户能否接受”这个根本问题上来。
Comments