一、什么是CQRS读模型的冷启动和全量投影重建
咱们先把复杂的概念拆成大白话讲。CQRS这个词,你可以简单理解成“把数据的读写操作分开做”——比如你写订单的时候,走一套专门管“写”的逻辑;别人查订单的时候,走另一套专门管“读”的逻辑。 读模型就是专门给查数据用的“数据副本”,比如你在电商网站上搜“近7天的热销商品”,后台查的不是原始的订单数据,而是提前整理好的热销商品读模型。 那冷启动是什么?就是读模型刚建起来或者坏了,里面啥数据都没有,得从原始数据(也就是“写模型”存的所有数据)里把读模型需要的内容一条一条转过来,这个过程就叫全量投影重建。说白了就是给空的读模型“填数据”。 举个最常见的例子:你做一个电商的“用户订单汇总表”读模型,用来快速查每个用户买过多少东西、花了多少钱。如果这个读模型的服务器硬盘坏了,数据全丢了,你就得从最开始的第一条订单开始,把所有用户的订单都重新转成这个汇总表的内容,这就是一次全量投影重建。
二、全量投影重建的实际场景
全量投影重建不是天天发生,但遇到下面这几种情况,你必须得做:
2.1 读模型故障
读模型一般是存在专门的数据库里,比如Redis、Elasticsearch或者专门的读库。如果这些数据库的硬盘坏了、删库了、或者数据被误改乱了,整个读模型的数据都废了,只能重新建。
2.2 读模型逻辑改了
比如原来的“用户订单汇总表”只统计近1年的订单,现在要改成统计所有历史订单;或者原来只统计实付金额,现在要加个“满减优惠总额”的字段。这时候旧的读模型内容不符合新逻辑,就得把所有数据重新转一遍,生成新的读模型。
2.3 新上线读模型
比如新做了一个“商品实时销量榜”的读模型,刚上线的时候里面啥都没有,得先把所有历史订单的数据转进去,才能给用户看。
三、全量投影重建的痛点:耗时和资源控制
全量投影重建最麻烦的两个问题,就是“太慢”和“太占资源”。 先说耗时:如果你的原始数据有1000万条订单,转成读模型的时候,每条都要做计算(比如算用户的总消费、商品的总销量),转完可能要几个小时甚至几天。这期间用户查读模型的数据要么是空的,要么是旧的,体验很差。 再说资源:转数据的时候,要从原始数据库(写模型)读数据,要做计算,还要把结果写到新的读模型里。这三步都占服务器的CPU、内存、硬盘读写、网络带宽。如果控制不好,会把写模型的数据库搞崩,导致正常的下单、支付操作都卡甚至失败,影响线上业务。
四、怎么解决耗时问题:让重建快起来
要缩短重建时间,核心就是“少做无用功”和“能快就快”,下面是具体的办法和例子。
4.1 分批次重建,不要一次性全转
不要把1000万条订单一次全读出来,分成100批,每批转10万条。这样做的好处是:第一,不会一次性占太多内存;第二,如果转的过程中出问题了,比如服务器断网了,只需要从最后一个转完的批次继续转,不用从头再来。 这里给个具体的代码例子,用Java写的,专门用来转订单数据到读模型,代码里加了详细注释。 技术栈:Java + MyBatis(用来读原始订单数据) + Redis(用来存读模型)
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.stereotype.Component;
import java.util.List;
@Component
public class OrderProjectionRebuilder {
@Autowired
private OrderMapper orderMapper; // MyBatis的Mapper,用来读原始订单数据
@Autowired
private StringRedisTemplate redisTemplate; // 用来存读模型的Redis客户端
// 重建读模型的主方法
public void rebuildReadModel(int batchSize) {
// 第一步:先查原始订单的总数量,用来算要分多少批
int totalOrderCount = orderMapper.getTotalOrderCount();
// 计算总批数:比如总订单1000万,每批10万,就是100批
int totalBatchCount = (int) Math.ceil((double) totalOrderCount / batchSize);
// 循环每一批,从第0批开始
for (int batchIndex = 0; batchIndex < totalBatchCount; batchIndex++) {
// 计算当前批次的起始位置:第0批从0开始,第1批从10万开始,以此类推
int offset = batchIndex * batchSize;
// 第二步:读当前批次的原始订单数据
List<Order> currentBatchOrders = orderMapper.getOrderBatch(offset, batchSize);
// 第三步:把当前批次的订单转成读模型需要的格式
for (Order order : currentBatchOrders) {
// 读模型的键:比如“user:123:order_summary”,值是该用户的订单汇总信息
String redisKey = "user:" + order.getUserId() + ":order_summary";
// 读模型的值:把订单的用户ID、总消费、总订单数拼起来(实际可以用JSON)
String summary = order.getUserId() + "," + order.getTotalAmount() + "," + order.getOrderCount();
// 第四步:把转好的数据写到读模型(Redis)里
redisTemplate.opsForValue().set(redisKey, summary);
}
// 打印进度,方便看转了多少
System.out.println("已完成第" + (batchIndex + 1) + "批,共" + totalBatchCount + "批");
}
}
}
这个代码的核心就是分批次,每批转完才转下一批,不会一次性占太多资源。
4.2 用并行处理,同时转多个批次
如果你的服务器有多核CPU,可以同时转多个批次,比如同时转2个批次,速度就能快一倍。但要注意不能开太多,不然CPU占满了反而慢。 比如上面的代码,你可以改成用线程池,同时跑2个批次的转任务:
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
// 重建读模型的主方法,加了并行处理
public void rebuildReadModelParallel(int batchSize) {
int totalOrderCount = orderMapper.getTotalOrderCount();
int totalBatchCount = (int) Math.ceil((double) totalOrderCount / batchSize);
// 创建线程池,固定2个线程,同时转2个批次
ExecutorService threadPool = Executors.newFixedThreadPool(2);
for (int batchIndex = 0; batchIndex < totalBatchCount; batchIndex++) {
int offset = batchIndex * batchSize;
// 把每个批次的转任务提交到线程池,同时跑
int finalBatchIndex = batchIndex;
threadPool.submit(() -> {
List<Order> currentBatchOrders = orderMapper.getOrderBatch(offset, batchSize);
for (Order order : currentBatchOrders) {
String redisKey = "user:" + order.getUserId() + ":order_summary";
String summary = order.getUserId() + "," + order.getTotalAmount() + "," + order.getOrderCount();
redisTemplate.opsForValue().set(redisKey, summary);
}
System.out.println("已完成第" + (finalBatchIndex + 1) + "批,共" + totalBatchCount + "批");
});
}
// 关闭线程池,等所有任务跑完
threadPool.shutdown();
}
这里要注意:线程池的大小不能超过CPU的核心数,比如你的服务器是4核的,最多开3个线程,留1个核给正常的业务用。
4.3 简化转的逻辑,减少计算量
如果转每条数据的时候要做很多复杂计算,比如要查用户的历史消费等级、要算商品的分类层级,那速度肯定慢。可以提前把这些需要用到的基础数据准备好,或者简化计算逻辑。 比如原来转订单的时候,要先查用户的会员等级,再算订单的折扣,再转成读模型。可以改成:在原始订单数据里就把会员等级、折扣都存好,转的时候直接用,不用再查。
五、怎么控制资源:别把线上业务搞崩
控制资源的核心是“限速”和“隔离”,不能让重建任务抢了正常业务的资源。
5.1 限制重建任务的CPU和内存
如果你的服务器是云服务器,或者用了Kubernetes(容器化部署),可以给重建任务单独设置CPU和内存的上限。比如给重建任务分配1核CPU、2G内存,就算它跑满了,也不会影响其他业务。 举个Kubernetes的配置例子,限制重建任务的资源:
apiVersion: batch/v1
kind: Job
metadata:
name: order-projection-rebuild
spec:
template:
spec:
containers:
- name: rebuild-container
image: your-rebuild-image:latest
# 限制CPU最多用1核,内存最多用2G
resources:
limits:
cpu: "1"
memory: "2Gi"
# 申请的资源,至少要1核、2G内存
requests:
cpu: "1"
memory: "2Gi"
restartPolicy: Never
这个配置的意思是:这个重建任务最多只能用1核CPU和2G内存,就算它想占更多,系统也会限制它。
5.2 限制读原始数据的速度
重建的时候要从原始数据库读数据,如果读的速度太快,会把原始数据库的硬盘读写占满,导致正常的下单操作(也要写原始数据库)卡。所以要限制读的速度,比如每秒钟最多读1000条订单。 怎么限制?可以在代码里加个等待时间,每转完一批,等1秒钟再转下一批。比如:
// 转完一批后,等1秒钟再转下一批
Thread.sleep(1000); // 单位是毫秒,1000毫秒就是1秒
你也可以根据业务情况调整等待时间,比如高峰的时候(比如晚上8点到10点)等2秒,低峰的时候(比如凌晨2点到4点)等0.5秒。
5.3 隔离读模型和写模型的数据库
最好的办法是把读模型的数据库和写模型的数据库分开,比如写模型用MySQL,读模型用Redis或者Elasticsearch。这样重建读模型的时候,就算把读模型的数据库占满了,也不会影响写模型的数据库,正常的下单操作还能正常进行。 比如你原来的架构是:用户下单→写MySQL(写模型);用户查订单→读MySQL(写模型)。改成:用户下单→写MySQL(写模型);转数据的时候→从MySQL读→转成读模型→写Redis(读模型);用户查订单→读Redis(读模型)。这样重建的时候只占Redis的资源,不影响MySQL。
六、重建过程中的注意事项
6.1 重建期间的读问题
重建的时候,读模型的数据是不完整的,用户查的时候可能会查到空数据或者旧数据。怎么解决?可以在重建期间给用户提示“数据正在更新,请稍后再试”,或者先让用户查旧的读模型,等重建完再切到新的。 比如你做一个“商品销量榜”的读模型,重建的时候,用户点“销量榜”,页面显示“销量榜正在更新,预计10分钟后完成”,这样用户体验不会太差。
6.2 重建期间的写问题
重建的时候,新的订单还在不断产生,这些新订单要不要转?如果不转,重建完的读模型就缺了这些新订单的数据。怎么解决?可以用“双写”的办法:新订单产生的时候,既写到写模型,也写到读模型;重建完之后,再把之前的新订单数据补到读模型里。 比如:你从10点开始重建读模型,转的是10点之前的订单;10点之后产生的新订单,每来一个就直接写到读模型里;11点重建完之后,读模型里既有10点之前的订单,也有10点之后的新订单,数据是完整的。
6.3 异常处理
重建的时候可能会出各种问题,比如网络断了、数据库连不上、代码出bug。所以要做异常处理,比如转某一批的时候出问题了,要记录下来,等问题解决了再单独转这一批,不用从头再来。 比如在代码里加try-catch:
try {
// 转当前批次的代码
} catch (Exception e) {
// 把出错的批次号、错误信息写到日志里
logger.error("转第" + batchIndex + "批出错,错误信息:" + e.getMessage());
// 可以选择暂停重建,或者跳过这一批后面再补
}
七、总结
全量投影重建是CQRS架构里一个很重要的操作,它的核心问题就是耗时太长和占资源太多。解决这些问题的办法主要有:分批次重建、并行处理、限制资源、隔离数据库、做好异常处理。 只要你把这些办法用对了,就能让重建又快又稳,不会影响线上业务。比如你原来转1000万条订单要10个小时,用了分批次加并行处理,可能2个小时就转完了;用了资源限制,也不会把线上的数据库搞崩。
Comments