CQRS 模式实战:用 Node.js 构建高性能读写分离应用

By | 2026年6月21日

为什么需要 CQRS?

在传统的 CRUD 应用中,我们通常使用同一个数据模型来处理读和写操作。然而,随着业务复杂度增加,这种模式会遇到几个问题:

  • 读写负载不均衡:查询操作往往比写入多得多,但共享同一数据源,导致写入性能受查询影响。
  • 数据模型不匹配:写入时需要满足业务规则(如唯一性校验),而读取时需要灵活的数据结构(如聚合统计),同一个模型难以兼顾。
  • 扩展性差:无法独立扩展读或写能力。

CQRS(Command Query Responsibility Segregation)模式将命令(写操作)和查询(读操作)分离,各自使用独立的数据模型和存储,从而解决上述问题。

实战场景:订单管理系统

我们将构建一个简单的订单管理系统,包含以下功能:

  • 创建订单(写)
  • 查询订单列表(读)
  • 查询订单详情(读)

我们将使用 Node.js、Express、MongoDB(写库)和 Redis(读库)来演示 CQRS。

项目结构


project/
├── src/
│   ├── command/          # 命令处理
│   │   ├── handlers/     # 命令处理器
│   │   └── models/       # 写模型(MongoDB)
│   ├── query/            # 查询处理
│   │   ├── handlers/     # 查询处理器
│   │   └── models/       # 读模型(Redis)
│   ├── events/           # 事件(可选,用于同步)
│   ├── routes/           # 路由
│   └── app.js
├── package.json
└── README.md

第一步:搭建项目基础


mkdir cqrs-demo && cd cqrs-demo
npm init -y
npm install express mongoose redis uuid

创建 src/app.js


const express = require('express');
const mongoose = require('mongoose');
const redis = require('redis');

const app = express();
app.use(express.json());

// 连接 MongoDB(写库)
mongoose.connect('mongodb://localhost:27017/cqrs_write', {
  useNewUrlParser: true,
  useUnifiedTopology: true
});

// 连接 Redis(读库)
const redisClient = redis.createClient();
redisClient.on('error', err => console.error('Redis error:', err));

// 将 Redis 客户端挂载到 app
app.set('redisClient', redisClient);

// 路由挂载(稍后实现)
app.use('/orders', require('./routes/orders'));

const PORT = process.env.PORT || 3000;
app.listen(PORT, () => {
  console.log(`Server running on port ${PORT}`);
});

第二步:定义写模型(MongoDB)

src/command/models/Order.js


const mongoose = require('mongoose');

const orderSchema = new mongoose.Schema({
  orderId: { type: String, required: true, unique: true },
  customerId: { type: String, required: true },
  items: [{
    productId: String,
    quantity: Number,
    price: Number
  }],
  totalAmount: { type: Number, required: true },
  status: { type: String, enum: ['pending', 'confirmed', 'shipped', 'delivered'], default: 'pending' },
  createdAt: { type: Date, default: Date.now },
  updatedAt: { type: Date, default: Date.now }
});

module.exports = mongoose.model('Order', orderSchema);

> 💡 写模型通常需要满足业务约束,比如订单总额必须等于各项之和。这里我们简单记录。

第三步:实现命令处理器

创建 src/command/handlers/createOrderHandler.js


const { v4: uuidv4 } = require('uuid');
const Order = require('../models/Order');

async function handleCreateOrder(command) {
  // 命令包含 customerId, items
  const { customerId, items } = command;

  // 计算总金额
  const totalAmount = items.reduce((sum, item) => sum + item.quantity * item.price, 0);

  const order = new Order({
    orderId: uuidv4(),
    customerId,
    items,
    totalAmount,
    status: 'pending',
    createdAt: new Date(),
    updatedAt: new Date()
  });

  await order.save();

  // 返回新创建的订单 ID
  return { orderId: order.orderId };
}

module.exports = { handleCreateOrder };

第四步:定义读模型(Redis)

读模型通常是为了优化查询性能,我们可以存储反范式化的数据。例如,在 Redis 中存储订单列表和详情。

创建 src/query/models/orderReadModel.js


const redisClient = require('../../app').get('redisClient'); // 注意:实际使用需注入

// 实际项目中,应该通过依赖注入或模块导出方式获取 redisClient。这里简化。
// 更好的做法:将 redisClient 作为参数传入。

const ORDER_LIST_KEY = 'orders:list';
const ORDER_DETAIL_KEY = 'order:';

async function addOrder(order) {
  // 存储订单详情(使用 Hash)
  await redisClient.hSet(`${ORDER_DETAIL_KEY}${order.orderId}`, {
    orderId: order.orderId,
    customerId: order.customerId,
    items: JSON.stringify(order.items),
    totalAmount: order.totalAmount,
    status: order.status,
    createdAt: order.createdAt.toISOString(),
    updatedAt: order.updatedAt.toISOString()
  });

  // 添加到有序集合,用于列表查询(按创建时间排序)
  await redisClient.zAdd(ORDER_LIST_KEY, {
    score: order.createdAt.getTime(),
    value: order.orderId
  });
}

async function getOrderList(page, pageSize) {
  const start = (page - 1) * pageSize;
  const end = start + pageSize - 1;
  // 获取指定范围内的订单 ID(按时间倒序)
  const orderIds = await redisClient.zRange(ORDER_LIST_KEY, start, end, { REV: true });
  if (orderIds.length === 0) return [];

  // 批量获取订单详情
  const pipeline = redisClient.multi();
  orderIds.forEach(id => {
    pipeline.hGetAll(`${ORDER_DETAIL_KEY}${id}`);
  });
  const details = await pipeline.exec();

  return details.map(d => ({
    ...d,
    items: JSON.parse(d.items || '[]'),
    totalAmount: parseFloat(d.totalAmount),
    createdAt: new Date(d.createdAt),
    updatedAt: new Date(d.updatedAt)
  }));
}

async function getOrderDetail(orderId) {
  const detail = await redisClient.hGetAll(`${ORDER_DETAIL_KEY}${orderId}`);
  if (!detail || !detail.orderId) return null;
  return {
    ...detail,
    items: JSON.parse(detail.items || '[]'),
    totalAmount: parseFloat(detail.totalAmount),
    createdAt: new Date(detail.createdAt),
    updatedAt: new Date(detail.updatedAt)
  };
}

module.exports = { addOrder, getOrderList, getOrderDetail };

> ⚠️ 注意:上面的代码直接引用了 ../../appredisClient,这是不好的实践。更好的方式是通过依赖注入或单独模块导出 Redis 客户端。这里为了简洁,我们假设在路由中注入。

第五步:实现查询处理器

创建 src/query/handlers/orderQueryHandler.js


const orderReadModel = require('../models/orderReadModel');

async function handleGetOrders(query) {
  const { page = 1, pageSize = 10 } = query;
  return await orderReadModel.getOrderList(page, pageSize);
}

async function handleGetOrderDetail(query) {
  const { orderId } = query;
  return await orderReadModel.getOrderDetail(orderId);
}

module.exports = { handleGetOrders, handleGetOrderDetail };

第六步:事件同步(写库 → 读库)

当命令执行成功后,我们需要将数据同步到读库。最简单的方式是在命令处理器中直接调用读模型的更新方法。但为了解耦,我们可以引入事件机制。

这里我们采用简单方式:在命令处理器中同步更新读库。

修改 src/command/handlers/createOrderHandler.js


const { v4: uuidv4 } = require('uuid');
const Order = require('../models/Order');
const orderReadModel = require('../../query/models/orderReadModel');

async function handleCreateOrder(command, redisClient) {
  const { customerId, items } = command;

  const totalAmount = items.reduce((sum, item) => sum + item.quantity * item.price, 0);

  const order = new Order({
    orderId: uuidv4(),
    customerId,
    items,
    totalAmount,
    status: 'pending',
    createdAt: new Date(),
    updatedAt: new Date()
  });

  await order.save();

  // 同步到读库
  await orderReadModel.addOrder(order, redisClient);

  return { orderId: order.orderId };
}

module.exports = { handleCreateOrder };

> 💡 注意:这里我们传递了 redisClient 参数,而不是直接引用。这样更灵活。

相应地,orderReadModel.addOrder 需要接受 redisClient


async function addOrder(order, redisClient) {
  await redisClient.hSet(`${ORDER_DETAIL_KEY}${order.orderId}`, {
    // ...
  });
  await redisClient.zAdd(ORDER_LIST_KEY, {
    score: order.createdAt.getTime(),
    value: order.orderId
  });
}

第七步:路由与控制器

创建 src/routes/orders.js


const express = require('express');
const router = express.Router();
const { handleCreateOrder } = require('../command/handlers/createOrderHandler');
const { handleGetOrders, handleGetOrderDetail } = require('../query/handlers/orderQueryHandler');

// 写操作:创建订单
router.post('/', async (req, res) => {
  try {
    const redisClient = req.app.get('redisClient');
    const result = await handleCreateOrder(req.body, redisClient);
    res.status(201).json(result);
  } catch (error) {
    res.status(500).json({ error: error.message });
  }
});

// 读操作:获取订单列表
router.get('/', async (req, res) => {
  try {
    const orders = await handleGetOrders(req.query);
    res.json(orders);
  } catch (error) {
    res.status(500).json({ error: error.message });
  }
});

// 读操作:获取订单详情
router.get('/:orderId', async (req, res) => {
  try {
    const order = await handleGetOrderDetail({ orderId: req.params.orderId });
    if (!order) {
      return res.status(404).json({ error: 'Order not found' });
    }
    res.json(order);
  } catch (error) {
    res.status(500).json({ error: error.message });
  }
});

module.exports = router;

第八步:测试运行

启动 MongoDB 和 Redis 后,运行 node src/app.js

测试写操作


curl -X POST http://localhost:3000/orders \
  -H "Content-Type: application/json" \
  -d '{"customerId": "cust123", "items": [{"productId": "prod1", "quantity": 2, "price": 10}, {"productId": "prod2", "quantity": 1, "price": 20}]}'

测试读操作


curl http://localhost:3000/orders
curl http://localhost:3000/orders/{orderId}

常见问题与最佳实践

1. 数据一致性

CQRS 带来了最终一致性。写库更新后,读库可能稍有延迟。对于需要强一致性的场景,可以考虑:

  • 使用事务性发件箱模式(Outbox Pattern)
  • 使用同步复制(如本例)

2. 事件溯源

更复杂的 CQRS 实现常与事件溯源(Event Sourcing)结合,将状态变更存储为事件序列,而不是当前状态。这提供了完整的审计日志和状态重建能力。

3. 扩展性

读写分离后,可以独立扩展读库(如增加 Redis 副本)和写库(如分片)。

4. 命令与查询的验证

命令处理器应该进行业务规则验证,而查询处理器只需返回数据。不要混在一起。

总结

本文通过一个订单管理系统的实战,演示了 CQRS 模式的基本实现:

  • 使用 MongoDB 作为写库,Redis 作为读库
  • 命令和查询分别由不同的处理器处理
  • 通过事件机制(或直接调用)同步数据

CQRS 适用于读写负载差异大、查询需求复杂的系统。但要注意,它增加了系统复杂度,不适合简单应用。

延伸阅读

如果你对事件溯源感兴趣,请关注我的下一篇文章:《事件溯源实战:用 Node.js 构建可追溯的订单系统》。