在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的分离优势。

事件溯源方案

  • 优点:最强的可追溯性和一致性;支持时间旅行和审计。
  • 缺点:学习成本高;查询需要投影,可能需要额外的专门读库;存储量大。

七、注意事项

  1. 不要为了用而用:如果你的系统读写量不大,或者读写模型一样,CQRS带来的复杂度可能大于收益。
  2. 事件总线的高可用:事件传递是整个系统的命脉,务必做好消息队列的集群、持久化和重试。
  3. 监控与补偿:建立一致性检查机制,比如定期对比写库和读库的数据差异,自动修复遗漏的事件。
  4. 领域事件设计:事件应该包含足够的信息,让读模型不需要再去查写库。例如订单创建事件应该包含所有需要展示的字段。
  5. 事务边界:在写模型中,发布事件和保存写库最好不要放在同一个本地事务里,因为如果事件发送失败,本地事务回滚会导致用户看到“操作成功”但实际没数据。更好的做法是先保存写库,然后发事件;如果事件发送失败,通过补偿任务重新发送。

八、文章总结

CQRS架构下的数据一致性问题,本质上是在“一致性”和“性能”之间做权衡。没有银弹,只能根据业务选择合适的策略。最终一致性适合大多数互联网应用,通过事件驱动的异步同步,可以同时获得高吞吐和较好的用户体验。对于核心的、强一致性的场景,可以考虑同步更新或事件溯源,但需要接受额外的复杂度。在实战中,一定要做好事件可靠传递、幂等处理和补偿机制。不要害怕不一致,但要对不一致窗口有清晰的认知,并告知业务方。记住,技术是为业务服务的,一致性方案的选择最终要回到“用户能否接受”这个根本问题上来。