一、先搞懂什么是CQRS(用买菜举例子)

很多人第一次听CQRS会觉得是高大上的架构,其实本质就是把“干事儿”和“查事儿”分开。举个大家天天接触的买菜例子: 你去菜市场买一斤青菜,先跟摊主说“要一斤青菜”(这是命令,是你要摊主执行的动作),摊主称完给你,你扫一眼确认斤数(这是查询,是你要获取的结果)。 在CQRS架构里,就把这俩活拆成两个独立的部分:

  • 命令端:专门处理“干事儿”的请求,比如下单、改地址、付款,核心是“改数据”,不关心改完的结果是什么。
  • 查询端:专门处理“查事儿”的请求,比如查订单状态、查地址、查余额,核心是“读数据”,不关心数据是怎么来的。 为啥要拆?举个实际的业务场景:比如电商的订单系统,一天可能有10万次下单(命令),但一天可能有100万次查订单(查询),分开后可以给查询端加缓存、加服务器,不用影响命令端的性能。

二、命令执行失败的坑:为什么CQRS会有数据不一致?

拆成两个端,就容易出“数据对不上”的问题。还是用买菜的例子: 你跟摊主说“要一斤青菜”(命令),摊主已经把青菜装袋了,但这时候你手机断网,付款没成功(命令执行失败),但摊主的记账本上已经记了“张三买了一斤青菜”(命令端的数据已经改了),你查自己的订单(查询端)却显示“付款失败”,这就不一致了。 再举个真实的业务场景:用户点了“修改收货地址”的按钮(命令),命令端已经把新地址写到自己的数据库里了,但更新查询端缓存的动作失败了,导致用户刷新页面(查询)看到的还是旧地址,这就是数据不一致。 CQRS里的命令执行失败,一般分两种情况:

2.1 命令本身执行失败

比如用户传的参数有问题(比如地址为空)、权限不够(不是自己的订单)、业务规则不允许(比如已经发货的订单不能改地址),这种情况命令端直接报错,不会改数据,问题不大。

2.2 命令执行成功,但后续动作失败

这才是最坑的!比如命令端已经把新地址写到自己的数据库了,但更新查询端缓存的动作失败了、或者同步到查询端数据库的动作失败了,这时候命令端的数据已经改了,查询端还是旧的,两边对不上。

三、核心方案:用事件兜底解决不一致

那怎么解决这个坑?核心思路是:命令端改完数据后,不直接去同步查询端,而是发一个“事件”,比如“地址修改了”的事件,查询端主动去消费这个事件,更新自己的数据。如果中间同步失败了,事件会留着,下次再消费。 啥是事件?还是用买菜的例子:摊主装完青菜,不是直接去改自己的记账本,而是先喊一声“张三买了一斤青菜”(发事件),记账的人(查询端)听到了再去改记账本。如果喊的时候记账的人没听到,摊主下次再喊一声(事件重试),直到记账的人听到为止。

3.1 事件兜底的完整流程

我们拿“修改收货地址”的业务来走一遍完整流程,所有动作都按顺序来:

  1. 用户点“修改地址”按钮,前端把新地址发给命令端。
  2. 命令端先验证参数(比如地址不能空),验证通过后,把新地址写到命令端的数据库(比如订单表的地址字段)。
  3. 命令端发一个“地址修改事件”,事件里包含订单ID、新地址、修改时间这些关键信息。
  4. 查询端一直在监听有没有“地址修改事件”,监听到后,把新地址更新到自己的数据库(或者缓存)。
  5. 如果第4步失败了(比如查询端数据库临时挂了),事件不会丢,会存在事件队列里,等查询端恢复了再消费。

3.2 事件队列的作用:为什么事件不会丢?

事件队列就像一个“消息盒子”,命令端发的事件先放到盒子里,查询端再从盒子里拿。如果查询端拿不到,事件就一直在盒子里,不会丢。 常见的事件队列有RabbitMQ、Kafka,我们举个简单的例子,用Kafka当事件队列: 首先明确技术栈:Node.js + Kafka(命令端) + Node.js + MongoDB(查询端)

3.2.1 命令端发事件的代码(Node.js + Kafka)

const { Kafka } = require('kafkajs');
const mongoose = require('mongoose'); // 命令端用MongoDB存订单

// 1. 连接Kafka
const kafka = new Kafka({
  clientId: 'order-command-service', // 命令端的唯一标识
  brokers: ['localhost:9092'] // Kafka的地址
});
const producer = kafka.producer(); // 生产者:用来发事件

// 2. 连接命令端的数据库(订单表)
mongoose.connect('mongodb://localhost:27017/order-command-db', { useNewUrlParser: true, useUnifiedTopology: true });
const Order = mongoose.model('Order', new mongoose.Schema({
  orderId: String,
  address: String,
  status: String
}));

// 3. 处理修改地址的命令
async function updateAddress(orderId, newAddress) {
  try {
    // 第一步:更新命令端的订单数据
    const updatedOrder = await Order.findOneAndUpdate(
      { orderId: orderId },
      { address: newAddress },
      { new: true } // 返回更新后的订单
    );
    if (!updatedOrder) throw new Error('订单不存在');

    // 第二步:发地址修改事件
    await producer.connect();
    await producer.send({
      topic: 'order-events', // 事件的主题,相当于消息盒子的名字
      messages: [
        {
          key: orderId, // 用订单ID当key,保证同一个订单的事件按顺序处理
          value: JSON.stringify({
            eventType: 'AddressUpdated', // 事件类型,告诉查询端是什么动作
            orderId: orderId,
            newAddress: newAddress,
            updateTime: new Date().toISOString()
          })
        }
      ]
    });
    console.log('地址修改成功,事件已发送');
  } catch (error) {
    console.error('修改地址失败:', error.message);
    // 这里要注意:如果发事件失败了,要回滚命令端的修改!
    // 因为命令端已经改了数据,但事件没发出去,查询端不会更新,会不一致
    if (error.message !== '订单不存在') {
      await Order.findOneAndUpdate(
        { orderId: orderId },
        { address: oldAddress } // 这里要存修改前的地址,用来回滚
      );
    }
  }
}

3.2.2 查询端消费事件的代码(Node.js + MongoDB)

const { Kafka } = require('kafkajs');
const mongoose = require('mongoose'); // 查询端用MongoDB存订单

// 1. 连接Kafka
const kafka = new Kafka({
  clientId: 'order-query-service', // 查询端的唯一标识
  brokers: ['localhost:9092']
});
const consumer = kafka.consumer({ groupId: 'order-query-group' }); // 消费者:用来拿事件

// 2. 连接查询端的数据库(订单表)
mongoose.connect('mongodb://localhost:27017/order-query-db', { useNewUrlParser: true, useUnifiedTopology: true });
const QueryOrder = mongoose.model('QueryOrder', new mongoose.Schema({
  orderId: String,
  address: String,
  status: String
}));

// 3. 消费事件的逻辑
async function consumeEvents() {
  await consumer.connect();
  await consumer.subscribe({ topic: 'order-events', fromBeginning: true }); // 从事件队列的开头开始拿,防止丢事件

  await consumer.run({
    eachMessage: async ({ topic, partition, message }) => {
      try {
        // 解析事件内容
        const event = JSON.parse(message.value.toString());
        console.log('收到事件:', event.eventType);

        // 根据事件类型处理
        if (event.eventType === 'AddressUpdated') {
          // 更新查询端的订单数据
          await QueryOrder.findOneAndUpdate(
            { orderId: event.orderId },
            { address: event.newAddress },
            { new: true }
          );
          console.log('查询端地址更新成功');
        }
      } catch (error) {
        console.error('处理事件失败:', error.message);
        // 处理失败的话,Kafka会自动重试,不用我们自己写逻辑
      }
    }
  });
}

// 启动消费
consumeEvents();

四、方案的细节:避免事件兜底出问题

事件兜底不是万能的,有几个细节要注意,不然反而会出更多问题。

4.1 事件要幂等:不能重复处理

啥是幂等?就是同一个事件处理1次和处理10次,结果是一样的。比如“地址修改”事件,如果查询端因为网络问题,两次收到同一个事件,处理两次后,地址不能变来变去。 怎么实现幂等?很简单,给每个事件加一个唯一的ID,查询端处理完一个事件后,把事件ID存下来,下次再收到同一个ID的事件,直接跳过。 比如修改事件的代码里加事件ID:

// 命令端发事件的时候加eventId
value: JSON.stringify({
  eventId: 'event_' + orderId + '_' + new Date().getTime(), // 唯一ID
  eventType: 'AddressUpdated',
  orderId: orderId,
  newAddress: newAddress,
  updateTime: new Date().toISOString()
})

查询端处理的时候,先查这个eventId有没有处理过:

// 查询端加一个事件处理表
const EventProcessed = mongoose.model('EventProcessed', new mongoose.Schema({
  eventId: String,
  processedAt: Date
}));

// 处理事件的时候
if (event.eventType === 'AddressUpdated') {
  // 先查事件有没有处理过
  const processed = await EventProcessed.findOne({ eventId: event.eventId });
  if (processed) {
    console.log('事件已处理,跳过');
    return;
  }
  // 处理事件
  await QueryOrder.findOneAndUpdate(...);
  // 存事件ID
  await EventProcessed.create({ eventId: event.eventId, processedAt: new Date() });
}

4.2 事件要按顺序处理

同一个订单的事件,必须按顺序处理。比如先改地址A,再改地址B,查询端必须先处理A的事件,再处理B的事件,不然地址会变成A,不对。 怎么保证顺序?Kafka里可以用“分区”,同一个订单ID的事件,会被分到同一个分区,查询端处理的时候,同一个分区的事件是按顺序来的。

4.3 死信队列:处理永久失败的事件

有些事件可能永远处理不了,比如事件里的订单ID不存在,或者新地址是非法的,这种事件会一直重试,占着事件队列的资源。 这时候就要用死信队列,把处理失败多次的事件,转到死信队列里,人工处理。比如Kafka里可以配置死信队列:

// 查询端配置死信队列
const dlq = kafka.consumer({ groupId: 'order-query-dlq-group' });
await dlq.connect();
await dlq.subscribe({ topic: 'order-events-dlq' });

// 处理事件失败的时候,转到死信队列
catch (error) {
  console.error('处理事件失败,转到死信队列:', error.message);
  await producer.send({
    topic: 'order-events-dlq',
    messages: [message]
  });
}

五、应用场景、优缺点和注意事项

5.1 应用场景

事件兜底适合哪些场景?

  • 对数据一致性要求高的业务:比如订单、支付、库存,一旦不一致会有资损。
  • 命令端和查询端分开的业务:比如电商、物流、金融,命令端和查询端的性能要求不一样。
  • 分布式系统:比如多个服务之间的数据同步,比如订单服务改了地址,要同步给物流服务、用户服务。

5.2 方案的优缺点

优点

  • 解决了CQRS的核心问题:数据不一致,用事件的方式,保证查询端最终和命令端一致。
  • 解耦:命令端和查询端不用直接依赖,命令端不用关心查询端有没有问题,查询端不用关心命令端怎么改数据。
  • 可扩展:事件队列可以扩展,查询端可以加多个,处理事件的速度更快。

缺点

  • 复杂度高:要加事件队列、要处理幂等、要处理死信队列,比直接同步复杂。
  • 最终一致:事件处理需要时间,所以查询端和命令端不是实时一致,是“最终一致”,比如用户改了地址,可能过1秒才在查询端看到。
  • 调试难:事件流是异步的,出问题了很难追踪,比如事件没发、没处理,要查Kafka、查日志,比同步难。

5.3 注意事项

  • 不要滥用:如果是小项目,命令端和查询端的性能要求不高,不用搞CQRS,更不用搞事件兜底,直接同步就好。
  • 监控事件:要监控事件的发送、消费情况,比如事件有没有积压、有没有失败,不然出问题了都不知道。
  • 测试边界:要测试各种失败场景,比如命令端发事件失败、查询端消费失败、事件重复,保证系统能正常处理。

六、总结

CQRS架构下命令执行失败导致的数据不一致,本质是命令端和查询端的同步问题,用事件兜底的方式,就是把同步变成异步,用事件队列保证事件不丢,让查询端主动同步数据。 整个流程的核心是:命令端改数据→发事件→查询端消费事件→更新自己的数据,中间的细节是幂等、顺序、死信队列,解决异步带来的问题。 这个方案适合对一致性要求高的分布式业务,但复杂度也高,要根据自己的项目情况选择。