为什么需要 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 };
> ⚠️ 注意:上面的代码直接引用了 ../../app 的 redisClient,这是不好的实践。更好的方式是通过依赖注入或单独模块导出 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 构建可追溯的订单系统》。