CQRS 模式实战:从理论到代码,构建高性能查询系统
引言
在传统的 CRUD 架构中,读和写操作共用同一数据模型,导致系统在高并发场景下性能瓶颈明显。例如,电商订单系统需要频繁查询订单状态,同时也要处理下单、支付等写操作。如果读写混用,复杂的查询会拖慢写操作的响应时间。
CQRS(Command Query Responsibility Segregation)模式通过分离命令(写)和查询(读)的职责,允许各自独立优化。本文将以一个简化版的电商订单系统为例,从零开始实现 CQRS,并结合事件溯源(Event Sourcing)来保证数据一致性。
项目准备
我们使用 Node.js + TypeScript 实现,数据库选用 PostgreSQL(写库)和 Redis(读库缓存)。首先初始化项目:
mkdir cqrs-demo
cd cqrs-demo
npm init -y
npm install typescript @types/node ts-node pg redis uuid
创建 tsconfig.json:
{
"compilerOptions": {
"target": "ES2020",
"module": "commonjs",
"outDir": "./dist",
"rootDir": "./src",
"strict": true,
"esModuleInterop": true
}
}
核心概念实现
1. 命令与命令处理器
命令(Command)代表一个写操作,如创建订单。命令处理器(Command Handler)负责执行命令并产生事件。
// src/commands/CreateOrderCommand.ts
export class CreateOrderCommand {
constructor(
public readonly orderId: string,
public readonly userId: string,
public readonly items: { productId: string; quantity: number }[],
public readonly totalAmount: number
) {}
}
// src/commands/CreateOrderHandler.ts
import { CreateOrderCommand } from './CreateOrderCommand';
import { EventStore } from '../events/EventStore';
import { OrderCreatedEvent } from '../events/OrderCreatedEvent';
export class CreateOrderHandler {
constructor(private eventStore: EventStore) {}
async handle(command: CreateOrderCommand): Promise<void> {
// 业务校验(简化)
if (!command.orderId || !command.userId) {
throw new Error('Invalid command');
}
// 生成事件
const event = new OrderCreatedEvent(
command.orderId,
command.userId,
command.items,
command.totalAmount,
new Date()
);
// 保存事件到事件存储
await this.eventStore.save(event);
}
}
2. 事件与事件存储
事件(Event)表示已经发生的事情,是不可变的。事件存储(Event Store)负责持久化事件,并支持按聚合ID查询。
// src/events/OrderCreatedEvent.ts
export class OrderCreatedEvent {
constructor(
public readonly orderId: string,
public readonly userId: string,
public readonly items: { productId: string; quantity: number }[],
public readonly totalAmount: number,
public readonly createdAt: Date
) {}
get eventType(): string {
return 'OrderCreated';
}
}
// src/events/EventStore.ts
import { Pool } from 'pg';
export class EventStore {
constructor(private pool: Pool) {}
async save(event: any): Promise<void> {
const query = `
INSERT INTO events (aggregate_id, event_type, event_data, created_at)
VALUES ($1, $2, $3, $4)
`;
const values = [
event.orderId,
event.eventType,
JSON.stringify(event),
event.createdAt
];
await this.pool.query(query, values);
}
async getEventsByAggregateId(aggregateId: string): Promise<any[]> {
const result = await this.pool.query(
'SELECT event_data FROM events WHERE aggregate_id = $1 ORDER BY created_at',
[aggregateId]
);
return result.rows.map(row => JSON.parse(row.event_data));
}
}
创建事件表:
CREATE TABLE events (
id SERIAL PRIMARY KEY,
aggregate_id VARCHAR(255) NOT NULL,
event_type VARCHAR(255) NOT NULL,
event_data JSONB NOT NULL,
created_at TIMESTAMP NOT NULL
);
CREATE INDEX idx_aggregate_id ON events (aggregate_id);
3. 读模型与投影
读模型(Read Model)是专为查询优化的数据视图。投影(Projection)负责消费事件并更新读模型。我们使用 Redis 作为缓存,同时维护一个 PostgreSQL 读库表。
// src/readmodels/OrderReadModel.ts
export interface OrderReadModel {
orderId: string;
userId: string;
items: { productId: string; quantity: number }[];
totalAmount: number;
status: string;
createdAt: Date;
}
// src/projections/OrderProjection.ts
import { EventStore } from '../events/EventStore';
import { Redis } from 'ioredis';
import { Pool } from 'pg';
export class OrderProjection {
constructor(
private eventStore: EventStore,
private redis: Redis,
private readDb: Pool
) {}
async project(event: any): Promise<void> {
if (event.eventType === 'OrderCreated') {
await this.handleOrderCreated(event);
}
}
private async handleOrderCreated(event: any): Promise<void> {
const readModel = {
orderId: event.orderId,
userId: event.userId,
items: event.items,
totalAmount: event.totalAmount,
status: 'created',
createdAt: event.createdAt
};
// 更新 Redis 缓存
await this.redis.set(
`order:${event.orderId}`,
JSON.stringify(readModel),
'EX',
3600
);
// 更新 PostgreSQL 读库
await this.readDb.query(
`INSERT INTO order_read (order_id, user_id, items, total_amount, status, created_at)
VALUES ($1, $2, $3, $4, $5, $6)
ON CONFLICT (order_id) DO UPDATE SET
items = EXCLUDED.items,
total_amount = EXCLUDED.total_amount,
status = EXCLUDED.status`,
[
readModel.orderId,
readModel.userId,
JSON.stringify(readModel.items),
readModel.totalAmount,
readModel.status,
readModel.createdAt
]
);
}
}
4. 查询处理器
查询处理器(Query Handler)负责处理读请求,直接从读模型返回数据。
// src/queries/GetOrderQuery.ts
export class GetOrderQuery {
constructor(public readonly orderId: string) {}
}
// src/queries/GetOrderHandler.ts
import { Redis } from 'ioredis';
import { Pool } from 'pg';
import { OrderReadModel } from '../readmodels/OrderReadModel';
export class GetOrderHandler {
constructor(
private redis: Redis,
private readDb: Pool
) {}
async handle(query: GetOrderQuery): Promise<OrderReadModel | null> {
// 先从缓存读取
const cached = await this.redis.get(`order:${query.orderId}`);
if (cached) {
return JSON.parse(cached);
}
// 缓存未命中,从读库读取
const result = await this.readDb.query(
'SELECT * FROM order_read WHERE order_id = $1',
[query.orderId]
);
if (result.rows.length === 0) {
return null;
}
const row = result.rows[0];
const readModel: OrderReadModel = {
orderId: row.order_id,
userId: row.user_id,
items: JSON.parse(row.items),
totalAmount: row.total_amount,
status: row.status,
createdAt: row.created_at
};
// 写入缓存
await this.redis.set(
`order:${query.orderId}`,
JSON.stringify(readModel),
'EX',
3600
);
return readModel;
}
}
组装与运行
1. 依赖注入与路由
使用简单的工厂模式组装组件:
// src/index.ts
import { Pool } from 'pg';
import Redis from 'ioredis';
import { EventStore } from './events/EventStore';
import { CreateOrderHandler } from './commands/CreateOrderHandler';
import { GetOrderHandler } from './queries/GetOrderHandler';
import { OrderProjection } from './projections/OrderProjection';
async function main() {
// 初始化数据库连接
const writeDb = new Pool({ connectionString: 'postgres://user:pass@localhost:5432/write_db' });
const readDb = new Pool({ connectionString: 'postgres://user:pass@localhost:5432/read_db' });
const redis = new Redis({ host: 'localhost', port: 6379 });
const eventStore = new EventStore(writeDb);
const createOrderHandler = new CreateOrderHandler(eventStore);
const getOrderHandler = new GetOrderHandler(redis, readDb);
const projection = new OrderProjection(eventStore, redis, readDb);
// 模拟事件处理(生产环境应使用消息队列)
// 实际项目中,事件存储后应发布到消息队列,由投影消费
// 这里简化:在命令处理器中直接调用投影
// 注意:更好的做法是异步解耦
// 示例:创建订单
const command = new CreateOrderCommand('order-123', 'user-456', [
{ productId: 'prod-1', quantity: 2 }
], 99.99);
await createOrderHandler.handle(command);
// 手动触发投影(实际应异步)
const events = await eventStore.getEventsByAggregateId('order-123');
for (const event of events) {
await projection.project(event);
}
// 查询订单
const query = new GetOrderQuery('order-123');
const order = await getOrderHandler.handle(query);
console.log('Order:', order);
}
main().catch(console.error);
> 注意:上述代码中,投影是同步调用的,仅用于演示。生产环境中,事件存储后应通过消息队列(如 RabbitMQ、Kafka)异步发布事件,由独立的投影服务消费,以避免写路径阻塞。
2. 数据库初始化脚本
创建读库表:
CREATE TABLE order_read (
order_id VARCHAR(255) PRIMARY KEY,
user_id VARCHAR(255) NOT NULL,
items JSONB NOT NULL,
total_amount DECIMAL(10,2) NOT NULL,
status VARCHAR(50) NOT NULL DEFAULT 'created',
created_at TIMESTAMP NOT NULL
);
进阶优化:事件溯源与聚合根
在 CQRS 中,聚合根(Aggregate Root)负责保证业务一致性。事件溯源(Event Sourcing)将聚合的状态存储为一系列事件,重建状态时回放事件。
实现一个简单的订单聚合根
// src/aggregates/OrderAggregate.ts
import { OrderCreatedEvent } from '../events/OrderCreatedEvent';
export class OrderAggregate {
public orderId: string = '';
public userId: string = '';
public status: string = '';
// 从事件重建状态
static fromEvents(events: any[]): OrderAggregate {
const aggregate = new OrderAggregate();
for (const event of events) {
aggregate.apply(event);
}
return aggregate;
}
apply(event: any): void {
if (event.eventType === 'OrderCreated') {
this.orderId = event.orderId;
this.userId = event.userId;
this.status = 'created';
}
// 其他事件处理...
}
}
这样,命令处理器可以从事件存储中获取所有事件,重建聚合状态,然后执行业务逻辑并产生新事件。
常见坑与最佳实践
- 最终一致性:CQRS 读模型是最终一致的,从写命令到读模型更新存在延迟。对于需要强一致的场景(如支付),应直接从写库读取或使用同步投影。
- 事件版本管理:事件结构可能变化,建议为事件添加版本号,并编写迁移脚本。
- 避免过度设计:并非所有系统都需要 CQRS。当你的系统读写负载不均衡,或查询复杂度远超写入时,才考虑引入。
- 消息队列可靠性:事件发布到消息队列时,需保证至少一次投递,并处理重复事件(幂等性)。
总结
本文通过一个电商订单系统,完整实现了 CQRS 模式,包括命令、事件、投影和查询。CQRS 将读写职责分离,使系统更灵活、可扩展。结合事件溯源,还能获得完整的审计日志。
下一步,你可以尝试:
- 引入消息队列(如 RabbitMQ)异步解耦命令和投影
- 实现更复杂的业务事件(如订单取消、支付成功)
- 部署到云环境,测试读写分离的性能提升
希望这篇文章能帮助你掌握 CQRS 的实战应用,构建更健壮的系统。