一、先搞懂啥是读写分离(CQRS的通俗版)
很多人一听到CQRS就觉得是高大上的架构,其实它本质就是把“读数据”和“写数据”这两件事拆开做,给它们各配一套专属的处理逻辑和数据存储。比如你做一个电商平台,用户下单(写操作)是一套流程,看商品列表(读操作)是另一套流程,互相不干扰。
举个生活里的例子就好懂了:小区的快递柜系统,快递员存快递(写操作)有自己的流程——扫快递码、开对应格子、关柜确认;用户取快递(读操作)是扫取件码、系统找格子、开柜。这两个操作的逻辑、权限、甚至用的系统模块都不一样,这就是CQRS的核心逻辑。
二、读写数据不一致的深层次原因
既然把读写拆开了,那肯定会有个问题:我刚写进去的数据,读的时候怎么找不到?或者读出来的是旧数据?很多人以为只是“延迟”,其实背后有好几个深层原因。
2.1 异步同步的“时间差”
为了让写操作更快,很多时候写数据不会直接同步到读的存储里,而是先写到一个临时的“写存储”,再通过消息队列(比如RabbitMQ、Kafka)慢慢同步到“读存储”。这个同步过程是有时间差的,短的几十毫秒,长的可能几秒甚至更久。
比如你刚在电商平台下了单(写操作),系统先把订单存到写存储,然后发了个“新订单”的消息给消息队列,读存储还没收到这个消息,你刷新订单列表(读操作)就看不到刚下的单,这就是异步同步的时间差导致的不一致。
2.2 同步失败的“断链”
同步过程不是100%可靠的,可能会出各种问题:消息队列里的消息丢了、读存储的同步服务挂了、网络断了……一旦同步断了,写进去的数据就永远同步不到读存储里,这就不是延迟的问题,是永久的不一致。
比如你写了个评论,写存储成功了,消息队列的消息因为网络问题丢了,读存储一直没收到“新评论”的消息,你刷新页面永远看不到自己刚写的评论,这就是同步失败导致的断链。
2.3 读写存储的“数据模型差异”
为了让读操作更快,读存储的数据模型会专门优化,比如把多个表的数据合并成一个大表(反范式设计),而写存储的模型是按业务逻辑拆分的(范式设计)。这种模型差异会导致同步的时候容易出问题,比如写存储的一个订单更新了,读存储的大表没更新对应的字段,就会出现数据不一致。
比如写存储里的订单表有“订单状态”“收货地址”两个字段,读存储的订单大表把这两个字段和商品表的“商品名称”“价格”合并成了一个表。如果写存储的订单状态更新了,同步的时候只更新了订单大表的“订单状态”,没更新关联的商品信息,就会出现读出来的商品价格是旧的,订单状态是新的这种不一致。
三、给大家一个完整的不一致场景示例
为了让大家更清楚,我们用一个具体的示例来模拟读写不一致的情况。这个示例用Spring Boot + RabbitMQ + MySQL做技术栈,实现一个简单的订单系统,故意做一些设计,让大家能看到不一致的效果。
3.1 技术栈说明
我们用的技术栈是:Spring Boot(后端框架)、RabbitMQ(消息队列)、MySQL(读写存储都用MySQL,但分开成写库和读库)。
3.2 写操作的代码
写操作的逻辑是:用户下单,先把订单存到写库,然后发一个“新订单”的消息到RabbitMQ。
// 写服务的下单方法
@Service
public class WriteOrderService {
// 写库的数据源
@Autowired
private JdbcTemplate writeJdbcTemplate;
// RabbitMQ的模板
@Autowired
private RabbitTemplate rabbitTemplate;
// 订单队列的名称
private static final String ORDER_QUEUE = "order.queue";
public void createOrder(Order order) {
// 1. 把订单存到写库
String sql = "INSERT INTO write_order (order_id, status, amount) VALUES (?, ?, ?)";
writeJdbcTemplate.update(sql, order.getOrderId(), order.getStatus(), order.getAmount());
// 2. 发消息到RabbitMQ,通知读库同步订单
rabbitTemplate.convertAndSend(ORDER_QUEUE, order);
System.out.println("写服务:订单" + order.getOrderId() + "已存写库,已发同步消息");
}
}
// 订单实体类
public class Order {
private String orderId;
private String status;
private double amount;
// 构造方法、getter、setter省略
}
3.3 读操作的代码
读操作的逻辑是:用户查询订单,直接从读库查。
// 读服务的查询订单方法
@Service
public class ReadOrderService {
// 读库的数据源
@Autowired
private JdbcTemplate readJdbcTemplate;
public Order getOrder(String orderId) {
// 1. 从读库查订单
String sql = "SELECT order_id, status, amount FROM read_order WHERE order_id = ?";
List<Order> orders = readJdbcTemplate.query(sql, new Object[]{orderId}, new BeanPropertyRowMapper<>(Order.class));
// 2. 返回订单,如果没查到返回null
return orders.isEmpty() ? null : orders.get(0);
}
}
3.4 同步服务的代码
同步服务的逻辑是:监听RabbitMQ的“新订单”消息,把订单存到读库。
// 同步服务的监听方法
@Component
public class OrderSyncService {
// 读库的数据源
@Autowired
private JdbcTemplate readJdbcTemplate;
// 订单队列的名称
private static final String ORDER_QUEUE = "order.queue";
// 监听RabbitMQ的订单消息
@RabbitListener(queues = ORDER_QUEUE)
public void syncOrder(Order order) {
// 1. 把订单存到读库
String sql = "INSERT INTO read_order (order_id, status, amount) VALUES (?, ?, ?)";
readJdbcTemplate.update(sql, order.getOrderId(), order.getStatus(), order.getAmount());
System.out.println("同步服务:订单" + order.getOrderId() + "已同步到读库");
}
}
3.5 模拟不一致的场景
现在我们来模拟一个场景:用户调用写服务创建一个订单,然后马上调用读服务查这个订单。
// 测试代码
public class Test {
@Autowired
private WriteOrderService writeOrderService;
@Autowired
private ReadOrderService readOrderService;
public void testInconsistency() {
// 1. 创建一个订单
Order order = new Order("123", "待支付", 99.9);
writeOrderService.createOrder(order);
// 2. 马上查这个订单
Order readOrder = readOrderService.getOrder("123");
// 3. 打印结果
if (readOrder == null) {
System.out.println("读服务:没查到订单,出现不一致!");
} else {
System.out.println("读服务:查到订单,数据一致!");
}
}
}
这个测试的结果大概率是“读服务:没查到订单,出现不一致!”,因为写服务刚把订单存到写库,发了消息,同步服务还没来得及处理消息,把订单存到读库,所以读服务查不到。
四、怎么保证最终一致?
最终一致的意思是:虽然短时间内读写数据可能不一致,但经过一段时间后,最终会变成一致的。我们可以用几个方法来保证这个最终一致。
4.1 写操作成功后,主动更新读缓存
如果我们用了缓存(比如Redis),写操作成功后,可以直接把数据写到缓存里,这样读操作就可以直接从缓存拿数据,不用等读存储同步。
比如上面的示例,写服务创建订单后,除了发消息到RabbitMQ,还可以把订单存到Redis:
// 写服务的下单方法修改后
@Service
public class WriteOrderService {
@Autowired
private JdbcTemplate writeJdbcTemplate;
@Autowired
private RabbitTemplate rabbitTemplate;
@Autowired
private StringRedisTemplate stringRedisTemplate; // 新增Redis模板
private static final String ORDER_QUEUE = "order.queue";
private static final String ORDER_KEY_PREFIX = "order:"; // Redis键的前缀
public void createOrder(Order order) {
// 1. 存写库
String sql = "INSERT INTO write_order (order_id, status, amount) VALUES (?, ?, ?)";
writeJdbcTemplate.update(sql, order.getOrderId(), order.getStatus(), order.getAmount());
// 2. 发消息到RabbitMQ
rabbitTemplate.convertAndSend(ORDER_QUEUE, order);
// 3. 主动把订单存到Redis,过期时间设为5分钟
String orderJson = JSON.toJSONString(order); // 用FastJSON把订单转成JSON
stringRedisTemplate.opsForValue().set(ORDER_KEY_PREFIX + order.getOrderId(), orderJson, 5, TimeUnit.MINUTES);
System.out.println("写服务:订单" + order.getOrderId() + "已存Redis");
}
}
然后读服务先从Redis查,查不到再从读库查:
// 读服务的查询方法修改后
@Service
public class ReadOrderService {
@Autowired
private JdbcTemplate readJdbcTemplate;
@Autowired
private StringRedisTemplate stringRedisTemplate;
private static final String ORDER_KEY_PREFIX = "order:";
public Order getOrder(String orderId) {
// 1. 先从Redis查
String orderJson = stringRedisTemplate.opsForValue().get(ORDER_KEY_PREFIX + orderId);
if (orderJson != null) {
return JSON.parseObject(orderJson, Order.class);
}
// 2. Redis查不到,再从读库查
String sql = "SELECT order_id, status, amount FROM read_order WHERE order_id = ?";
List<Order> orders = readJdbcTemplate.query(sql, new Object[]{orderId}, new BeanPropertyRowMapper<>(Order.class));
return orders.isEmpty() ? null : orders.get(0);
}
}
这样修改后,写操作成功后马上把数据写到Redis,读操作先从Redis拿,就不会出现短时间内查不到的情况,保证了数据的一致性。
4.2 重试机制:解决同步失败的问题
为了防止同步失败,我们可以给同步服务加重试机制:如果同步失败(比如读库挂了、网络断了),就把消息重新发回RabbitMQ,过一段时间再重试。
RabbitMQ本身就支持死信队列(DLX),可以用来做重试。我们可以配置一个死信队列,当同步服务处理消息失败时,消息会被发到死信队列,然后过一段时间再重新发送到原队列,让同步服务再处理一次。
比如我们给同步服务加重试的配置:
// 死信队列的配置
@Configuration
public class RabbitMqConfig {
// 订单队列的名称
private static final String ORDER_QUEUE = "order.queue";
// 死信交换机的名称
private static final String DLX_EXCHANGE = "dlx.exchange";
// 死信队列的名称
private static final String DLX_QUEUE = "dlx.queue";
// 死信路由键的名称
private static final String DLX_ROUTING_KEY = "dlx.routing.key";
// 配置订单队列,指定死信交换机
@Bean
public Queue orderQueue() {
return QueueBuilder.durable(ORDER_QUEUE)
.withArgument("x-dead-letter-exchange", DLX_EXCHANGE) // 指定死信交换机
.withArgument("x-dead-letter-routing-key", DLX_ROUTING_KEY) // 指定死信路由键
.build();
}
// 配置死信交换机
@Bean
public DirectExchange dlxExchange() {
return new DirectExchange(DLX_EXCHANGE);
}
// 配置死信队列
@Bean
public Queue dlxQueue() {
return QueueBuilder.durable(DLX_QUEUE)
.withArgument("x-message-ttl", 3000) // 消息在死信队列里的过期时间,3秒后重新发送
.build();
}
// 绑定死信交换机和死信队列
@Bean
public Binding dlxBinding() {
return BindingBuilder.bind(dlxQueue()).to(dlxExchange()).with(DLX_ROUTING_KEY);
}
}
然后同步服务处理消息失败时,抛出异常,RabbitMQ会自动把消息发到死信队列:
// 同步服务的监听方法修改后
@Component
public class OrderSyncService {
@Autowired
private JdbcTemplate readJdbcTemplate;
private static final String ORDER_QUEUE = "order.queue";
@RabbitListener(queues = ORDER_QUEUE)
public void syncOrder(Order order) {
try {
// 1. 把订单存到读库
String sql = "INSERT INTO read_order (order_id, status, amount) VALUES (?, ?, ?)";
readJdbcTemplate.update(sql, order.getOrderId(), order.getStatus(), order.getAmount());
System.out.println("同步服务:订单" + order.getOrderId() + "已同步到读库");
} catch (Exception e) {
// 2. 同步失败,抛出异常,消息会被发到死信队列
System.out.println("同步服务:订单" + order.getOrderId() + "同步失败,将重试");
throw new RuntimeException("同步订单失败", e);
}
}
}
这样就算同步服务第一次处理消息失败,3秒后消息会重新发送到订单队列,同步服务会再处理一次,解决了同步失败的问题。
4.3 对账机制:兜底保证最终一致
就算有了重试机制,还是有可能出现同步失败的情况(比如读库永久挂了),这时候我们需要一个兜底的对账机制:定期把写库的数据和读库的数据做对比,找出不一致的地方,然后修正。
比如我们可以每天凌晨运行一个定时任务,对比写库和读库的订单数据:
// 对账服务的定时任务
@Component
public class OrderReconciliationService {
@Autowired
private JdbcTemplate writeJdbcTemplate;
@Autowired
private JdbcTemplate readJdbcTemplate;
// 每天凌晨1点执行对账
@Scheduled(cron = "0 0 1 * * ?")
public void reconcileOrders() {
System.out.println("开始订单对账");
// 1. 从写库查所有订单
String writeSql = "SELECT order_id, status, amount FROM write_order";
List<Order> writeOrders = writeJdbcTemplate.query(writeSql, new BeanPropertyRowMapper<>(Order.class));
// 2. 从读库查所有订单
String readSql = "SELECT order_id, status, amount FROM read_order";
List<Order> readOrders = readJdbcTemplate.query(readSql, new BeanPropertyRowMapper<>(Order.class));
// 3. 把读库的订单转成Map,方便对比
Map<String, Order> readOrderMap = readOrders.stream()
.collect(Collectors.toMap(Order::getOrderId, Function.identity()));
// 4. 对比写库和读库的订单
for (Order writeOrder : writeOrders) {
Order readOrder = readOrderMap.get(writeOrder.getOrderId());
// 5. 如果读库没有这个订单,或者订单数据不一致,就修正
if (readOrder == null || !writeOrder.equals(readOrder)) {
System.out.println("发现不一致的订单:" + writeOrder.getOrderId() + ",开始修正");
// 把写库的订单同步到读库
String updateSql = "INSERT INTO read_order (order_id, status, amount) VALUES (?, ?, ?) ON DUPLICATE KEY UPDATE status = ?, amount = ?";
readJdbcTemplate.update(updateSql, writeOrder.getOrderId(), writeOrder.getStatus(), writeOrder.getAmount(), writeOrder.getStatus(), writeOrder.getAmount());
}
}
System.out.println("订单对账完成");
}
}
这个定时任务会每天对比写库和读库的订单,找出不一致的地方,然后把写库的数据同步到读库,兜底保证最终一致。
五、应用场景、优缺点和注意事项
5.1 应用场景
CQRS架构适合什么场景呢?
- 读写比例差异大的场景:比如电商平台,读操作(看商品、看订单)远多于写操作(下单、改地址),这时候把读写拆开,专门优化读操作的性能,能大幅提升系统的整体性能。
- 业务逻辑复杂的场景:比如金融系统,写操作(转账、存款)的逻辑很复杂,需要严格的权限控制和事务处理,读操作(查余额、查交易记录)的逻辑很简单,这时候把读写拆开,各配一套逻辑,能让系统更清晰。
- 数据模型差异大的场景:比如内容管理系统(CMS),写操作(发布文章)需要按内容的结构来存数据,读操作(看文章列表)需要按时间、分类来聚合数据,这时候把读写拆开,专门优化读数据的模型,能让读操作更快。
5.2 技术优缺点
优点
- 性能提升:读写分开后,可以各自优化,比如读操作可以用缓存、反范式的存储,写操作可以用高性能的事务存储,提升整体性能。
- 架构清晰:读写的逻辑分开,各自的职责更清晰,开发和维护更方便,比如写操作的团队只需要关注写逻辑,读操作的团队只需要关注读逻辑。
- 可扩展性强:读写可以各自扩展,比如读操作压力大的时候,只需要增加读存储的节点,不用动写存储的节点,节省成本。
缺点
- 架构复杂:比传统的单存储架构多了读写分离、消息队列、同步服务、对账机制等模块,架构更复杂,开发和维护的成本更高。
- 数据不一致风险:因为读写是异步同步的,所以会有数据不一致的风险,需要额外的机制来保证最终一致,增加了开发的难度。
- 调试困难:出现数据不一致的问题时,需要同时查写存储、读存储、消息队列、同步服务等多个模块,调试起来更困难。
5.3 注意事项
- 不要盲目用CQRS:只有当你的系统真的有读写比例差异大、业务逻辑复杂、数据模型差异大等需求时,才考虑用CQRS,不要为了用而用,增加不必要的复杂度。
- 保证同步的可靠性:异步同步是CQRS的核心,一定要保证同步的可靠性,比如用重试机制、死信队列、对账机制等,减少数据不一致的风险。
- 做好监控:一定要监控写存储、读存储、消息队列、同步服务的状态,比如同步的延迟、同步失败的次数、对账的结果等,及时发现问题。
- 简化读操作的逻辑:读操作的逻辑一定要简单,尽量避免复杂的业务逻辑,因为读操作的压力大,复杂的逻辑会影响读操作的性能。
六、总结
CQRS架构本质就是把读写操作拆开,各自优化,提升系统的性能和可扩展性。但读写拆开后,会因为异步同步的时间差、同步失败、数据模型差异等原因导致数据不一致。要保证最终一致,我们可以用主动更新读缓存、重试机制、对账机制等方法。
在使用CQRS架构时,一定要结合自己的业务需求,不要盲目用,同时要做好同步的可靠性、监控等工作,减少数据不一致的风险。
评论
围绕“CQRS架构中读写模型数据不一致的深层次原因与最终一致性保障方案”参与讨论