CQRS 模式实战:用 TypeScript 构建可扩展的读写分离架构
摘要
CQRS(Command Query Responsibility Segregation)是一种将读操作和写操作分离的架构模式,能显著提升系统的可扩展性和维护性。本文将通过一个电商订单系统的完整案例,手把手教你实现 CQRS 模式,涵盖命令与查询分离、事件溯源、读写模型独立优化等核心实践,并分享生产环境中的常见陷阱与解决方案。
为什么需要 CQRS?
在传统 CRUD 应用中,我们通常对同一数据模型进行读写操作。随着业务复杂度的提升,这种模式会遇到瓶颈:
- 读写负载不均衡:高频查询和低频写入共用同一个模型,导致资源争抢。
- 查询性能瓶颈:复杂的报表查询需要多表 JOIN,而写操作要求数据一致性,两者难以兼顾。
- 团队协作困难:读写逻辑耦合在同一个服务中,不同团队修改同一代码容易冲突。
CQRS 通过将命令(写)和查询(读)分离到不同的模型和服务中,解决了上述问题。
案例:电商订单系统
我们将构建一个简化版的电商订单系统,支持以下功能:
- 创建订单:用户提交购物车生成订单。
- 查询订单:按用户 ID 查询订单列表。
- 查询订单统计:获取每日订单数量。
我们将采用 CQRS 模式,将写操作(命令)和读操作(查询)分离。
技术栈
- TypeScript
- Express.js
- MongoDB(写模型)
- Redis(读模型缓存)
- Node.js
项目结构
src/
├── commands/ # 命令处理
│ ├── CreateOrderCommand.ts
│ └── CreateOrderHandler.ts
├── queries/ # 查询处理
│ ├── GetOrdersQuery.ts
│ └── GetOrdersHandler.ts
├── models/ # 数据模型
│ ├── Order.ts # 写模型(Mongoose)
│ └── OrderReadModel.ts # 读模型(Redis)
├── events/ # 事件
│ └── OrderCreatedEvent.ts
├── infrastructure/ # 基础设施
│ ├── EventBus.ts
│ └── Database.ts
└── index.ts # 入口
第一步:定义命令和查询
命令和查询是 CQRS 的核心概念。命令表示一个意图(如“创建订单”),通常以动词开头;查询表示一个请求(如“获取订单列表”)。
创建订单命令
// src/commands/CreateOrderCommand.ts
export interface CreateOrderCommand {
userId: string;
items: Array<{ productId: string; quantity: number; price: number }>;
shippingAddress: string;
}
获取订单查询
// src/queries/GetOrdersQuery.ts
export interface GetOrdersQuery {
userId: string;
page: number;
pageSize: number;
}
第二步:实现命令处理程序
命令处理程序负责验证业务规则并更新写模型。
// src/commands/CreateOrderHandler.ts
import { CreateOrderCommand } from './CreateOrderCommand';
import Order from '../models/Order';
import { EventBus } from '../infrastructure/EventBus';
import { OrderCreatedEvent } from '../events/OrderCreatedEvent';
export class CreateOrderHandler {
async handle(command: CreateOrderCommand): Promise<string> {
// 1. 验证业务规则
if (!command.items || command.items.length === 0) {
throw new Error('订单必须包含至少一个商品');
}
// 2. 创建订单实体(写模型)
const order = new Order({
userId: command.userId,
items: command.items,
shippingAddress: command.shippingAddress,
totalAmount: command.items.reduce((sum, item) => sum + item.price * item.quantity, 0),
status: 'pending',
createdAt: new Date(),
});
// 3. 保存到 MongoDB
await order.save();
// 4. 发布事件
const event: OrderCreatedEvent = {
orderId: order._id.toString(),
userId: order.userId,
totalAmount: order.totalAmount,
createdAt: order.createdAt,
};
EventBus.publish(event);
return order._id.toString();
}
}
注意:
> 命令处理程序不应返回任何数据(除了标识符),这是 CQRS 的重要原则。如果需要确认结果,可以返回订单 ID 或抛出异常。
第三步:实现查询处理程序
查询处理程序从读模型(Redis)或经过优化的数据库中读取数据。
// src/queries/GetOrdersHandler.ts
import { GetOrdersQuery } from './GetOrdersQuery';
import { OrderReadModel } from '../models/OrderReadModel';
export class GetOrdersHandler {
async handle(query: GetOrdersQuery): Promise<any> {
// 从 Redis 读取缓存数据
const cacheKey = `orders:${query.userId}:page:${query.page}`;
const cached = await OrderReadModel.get(cacheKey);
if (cached) {
return JSON.parse(cached);
}
// 如果缓存未命中,从 MongoDB 读取(但使用专门优化的集合)
const orders = await OrderReadModel.findByUserId(query.userId, query.page, query.pageSize);
// 写入缓存
await OrderReadModel.set(cacheKey, JSON.stringify(orders), 'EX', 60);
return orders;
}
}
第四步:事件驱动同步
当写模型更新后,通过事件同步到读模型。这里我们使用简单的内存事件总线。
// src/events/OrderCreatedEvent.ts
export interface OrderCreatedEvent {
orderId: string;
userId: string;
totalAmount: number;
createdAt: Date;
}
// src/infrastructure/EventBus.ts
type EventHandler = (event: any) => void;
export class EventBus {
private static handlers: Map<string, EventHandler[]> = new Map();
static subscribe(eventType: string, handler: EventHandler) {
if (!this.handlers.has(eventType)) {
this.handlers.set(eventType, []);
}
this.handlers.get(eventType)!.push(handler);
}
static publish(event: any) {
const eventType = event.constructor.name;
const handlers = this.handlers.get(eventType) || [];
handlers.forEach(handler => handler(event));
}
}
同步读模型的事件处理
// 在应用启动时注册
import { EventBus } from './infrastructure/EventBus';
import { OrderCreatedEvent } from './events/OrderCreatedEvent';
import { OrderReadModel } from './models/OrderReadModel';
EventBus.subscribe('OrderCreatedEvent', async (event: OrderCreatedEvent) => {
// 更新读模型(例如:更新 Redis 缓存或专门用于查询的 MongoDB 集合)
await OrderReadModel.updateOnOrderCreated(event);
});
第五步:定义数据模型
写模型(MongoDB Schema)
// src/models/Order.ts
import mongoose from 'mongoose';
const orderSchema = new mongoose.Schema({
userId: { type: String, required: true },
items: [{
productId: String,
quantity: Number,
price: Number,
}],
shippingAddress: String,
totalAmount: Number,
status: { type: String, default: 'pending' },
createdAt: { type: Date, default: Date.now },
});
export default mongoose.model('Order', orderSchema);
读模型(Redis 操作封装)
// src/models/OrderReadModel.ts
import redis from 'redis';
import { promisify } from 'util';
const client = redis.createClient();
const getAsync = promisify(client.get).bind(client);
const setAsync = promisify(client.set).bind(client);
export class OrderReadModel {
static async get(key: string): Promise<string | null> {
return getAsync(key);
}
static async set(key: string, value: string, mode: string, duration: number) {
return setAsync(key, value, mode, duration);
}
static async findByUserId(userId: string, page: number, pageSize: number) {
// 实际项目中可以从一个专门用于查询的 MongoDB 集合中读取
// 这里简化处理
return [];
}
static async updateOnOrderCreated(event: any) {
// 更新缓存或数据库
// 例如:清除相关用户的缓存
const cacheKey = `orders:${event.userId}:*`;
// 使用 Redis SCAN 清除模式匹配的键
}
}
第六步:连接路由
在 Express 中,我们将命令和查询路由到对应的处理程序。
// src/index.ts
import express from 'express';
import { CreateOrderHandler } from './commands/CreateOrderHandler';
import { GetOrdersHandler } from './queries/GetOrdersHandler';
const app = express();
app.use(express.json());
// 写操作路由
app.post('/orders', async (req, res) => {
try {
const handler = new CreateOrderHandler();
const orderId = await handler.handle(req.body);
res.status(201).json({ orderId });
} catch (error) {
res.status(400).json({ error: error.message });
}
});
// 读操作路由
app.get('/orders', async (req, res) => {
try {
const handler = new GetOrdersHandler();
const orders = await handler.handle({
userId: req.query.userId as string,
page: parseInt(req.query.page as string) || 1,
pageSize: parseInt(req.query.pageSize as string) || 10,
});
res.json(orders);
} catch (error) {
res.status(500).json({ error: error.message });
}
});
app.listen(3000, () => console.log('Server running on port 3000'));
生产环境中的坑与最佳实践
1. 最终一致性
CQRS 通常与事件溯源结合,但事件同步存在延迟。在电商系统中,用户下单后可能需要等待几毫秒才能查询到订单。解决方法:
- 在写操作完成后,立即将数据写入读模型(同步更新),但会牺牲解耦性。
- 使用 Saga 模式确保跨服务的事务一致性。
2. 命令验证与错误处理
命令处理程序应返回统一错误格式,避免泄露内部细节。
3. 查询优化
读模型可以根据查询需求进行反范式化,例如将用户信息和订单信息合并存储,减少 JOIN。
4. 事件版本管理
当事件结构变化时,需要处理版本兼容。可以使用 Avro 或 Protocol Buffers 进行序列化。
总结
本文通过一个电商订单系统案例,演示了如何使用 TypeScript 实现 CQRS 模式。核心要点:
- 命令和查询分离,各自拥有独立的处理逻辑和数据模型。
- 写模型使用 MongoDB 保证数据一致性,读模型使用 Redis 提升查询性能。
- 通过事件总线实现写模型到读模型的异步同步。
CQRS 并非银弹,对于简单 CRUD 应用可能过度设计。但在需要高可扩展性、读写负载差异大的场景下,它能带来显著的收益。
延伸阅读
- 事件溯源(Event Sourcing)与 CQRS 的结合
- Axon Framework(Java 的 CQRS 框架)
- 使用 Kafka 作为事件总线实现分布式 CQRS
如果你在项目中实践过 CQRS,欢迎在评论区分享你的经验!