一、什么是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个小时就转完了;用了资源限制,也不会把线上的数据库搞崩。