CQRS 模式实战:从理论到代码,构建高性能查询系统

By | 2026年6月15日

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';
    }
    // 其他事件处理...
  }
}

这样,命令处理器可以从事件存储中获取所有事件,重建聚合状态,然后执行业务逻辑并产生新事件。

常见坑与最佳实践

  1. 最终一致性:CQRS 读模型是最终一致的,从写命令到读模型更新存在延迟。对于需要强一致的场景(如支付),应直接从写库读取或使用同步投影。
  2. 事件版本管理:事件结构可能变化,建议为事件添加版本号,并编写迁移脚本。
  3. 避免过度设计:并非所有系统都需要 CQRS。当你的系统读写负载不均衡,或查询复杂度远超写入时,才考虑引入。
  4. 消息队列可靠性:事件发布到消息队列时,需保证至少一次投递,并处理重复事件(幂等性)。

总结

本文通过一个电商订单系统,完整实现了 CQRS 模式,包括命令、事件、投影和查询。CQRS 将读写职责分离,使系统更灵活、可扩展。结合事件溯源,还能获得完整的审计日志。

下一步,你可以尝试:

  • 引入消息队列(如 RabbitMQ)异步解耦命令和投影
  • 实现更复杂的业务事件(如订单取消、支付成功)
  • 部署到云环境,测试读写分离的性能提升

希望这篇文章能帮助你掌握 CQRS 的实战应用,构建更健壮的系统。