事件溯源与CQRS架构实战:订单系统从CRUD到事件驱动的设计演进

订单系统为何需要从CRUD演进到事件溯源

事件溯源(Event Sourcing)与CQRS(Command Query Responsibility Segregation)常被一起讨论,但它们解决的是不同维度的问题。事件溯源解决的是状态可审计性——不存当前状态,而是存导致状态变化的所有事件序列,任何时刻的状态都可从事件流重放得到。CQRS解决的是读写模型分离——写模型优化事务一致性,读模型优化查询性能。订单系统天然需要两者:订单状态变更需要完整审计轨迹(合规要求),读写负载特征差异大(写操作重一致性,读操作重聚合查询)。

事件溯源核心模型设计

传统CRUD模式下,订单状态直接写入orders表。事件溯源模式下,orders表不存在,取而代之的是order_events事件流:

-- 事件存储表
CREATE TABLE order_events (
    event_id UUID PRIMARY KEY,
    aggregate_id UUID NOT NULL,
    event_type VARCHAR(100) NOT NULL,
    event_data JSONB NOT NULL,
    version INT NOT NULL,
    created_at TIMESTAMP DEFAULT NOW(),
    UNIQUE(aggregate_id, version)
);

CREATE INDEX idx_order_events_aggregate
  ON order_events(aggregate_id, version);

订单生命周期对应的事件类型:

public class OrderEvent {
    private String eventType;
    private String aggregateId;
    private int version;
    private Map<String, Object> data;
    private LocalDateTime createdAt;
}

// 事件类型定义
// OrderCreated       - 订单创建
// ItemAdded          - 添加商品
// ItemRemoved        - 移除商品
// OrderConfirmed     - 确认订单
// PaymentReceived    - 收到付款
// OrderShipped       - 订单发货
// OrderDelivered     - 订单送达
// OrderCancelled     - 订单取消

聚合根(Aggregate)是事件溯源的核心概念,封装业务规则和状态转换逻辑:

public class OrderAggregate {
    private UUID orderId;
    private List<OrderItem> items;
    private OrderStatus status;
    private int version;

    // 从事件流重建状态
    public static OrderAggregate fromEvents(
            List<OrderEvent> events) {
        OrderAggregate order = new OrderAggregate();
        for (OrderEvent event : events) {
            order.apply(event);
        }
        return order;
    }

    private void apply(OrderEvent event) {
        switch (event.getEventType()) {
            case 'OrderCreated':
                this.orderId = UUID.fromString(
                    event.getData().get('orderId').toString());
                this.status = OrderStatus.CREATED;
                break;
            case 'ItemAdded':
                this.items.add(new OrderItem(
                    event.getData().get('sku').toString(),
                    Integer.parseInt(
                        event.getData().get('quantity').toString())
                ));
                break;
            case 'OrderConfirmed':
                this.status = OrderStatus.CONFIRMED;
                break;
            case 'PaymentReceived':
                this.status = OrderStatus.PAID;
                break;
        }
        this.version = event.getVersion();
    }

    // 业务命令 - 校验 + 产生事件
    public OrderEvent confirm() {
        if (this.status != OrderStatus.CREATED) {
            throw new IllegalStateException(
                '只有CREATED状态可确认');
        }
        if (this.items.isEmpty()) {
            throw new IllegalStateException(
                '空订单不可确认');
        }
        return OrderEvent.create(
            'OrderConfirmed', this.orderId,
            this.version + 1, Map.of());
    }
}

CQRS读写分离模型实现

写模型(Command端)直接操作事件流,读模型(Query端)通过事件投影(Projection)构建物化视图。两者数据源完全独立,通过事件总线同步:

// 写模型 - 事件存储
@Service
public class OrderCommandService {

    @Transactional
    public UUID createOrder(CreateOrderCommand cmd) {
        UUID orderId = UUID.randomUUID();
        OrderEvent event = OrderEvent.create(
            'OrderCreated', orderId, 1,
            Map.of('customerId',
                cmd.getCustomerId().toString())
        );
        eventRepository.append(event);
        eventBus.publish(event);
        return orderId;
    }
}
-- 读模型 - 物化视图
CREATE TABLE order_views (
    order_id UUID PRIMARY KEY,
    customer_id UUID,
    total_amount DECIMAL(12,2),
    status VARCHAR(20),
    item_count INT,
    updated_at TIMESTAMP
);
// 事件投影处理器 - 订阅事件更新读模型
@Component
public class OrderViewProjector {

    @EventHandler
    public void handle(OrderCreated event) {
        jdbcTemplate.update(
            'INSERT INTO order_views ' +
            '(order_id, customer_id, status, ' +
            'item_count, total_amount) ' +
            'VALUES (?, ?, \'CREATED\', 0, 0)',
            event.getAggregateId(),
            event.getData().get('customerId')
        );
    }

    @EventHandler
    public void handle(ItemAdded event) {
        jdbcTemplate.update(
            'UPDATE order_views SET ' +
            'item_count = item_count + 1, ' +
            'total_amount = total_amount + ? ' +
            'WHERE order_id = ?',
            event.getData().get('price'),
            event.getAggregateId()
        );
    }
}

快照机制优化长事件流性能

事件溯源的固有代价是状态重建开销。一个存在两年的订单如果经历了数百次状态变更,每次读取都要回放全部事件,性能不可接受。快照机制是标准解决方案:每隔N个事件保存一次聚合根的完整状态快照,重建时从最近快照开始回放。

CREATE TABLE order_snapshots (
    aggregate_id UUID,
    version INT,
    snapshot_data JSONB NOT NULL,
    created_at TIMESTAMP DEFAULT NOW(),
    PRIMARY KEY (aggregate_id, version)
);
public class SnapshotRepository {
    private static final int SNAPSHOT_INTERVAL = 50;

    public Optional<Snapshot> loadLatest(
            UUID aggregateId) {
        return jdbcTemplate.query(
            'SELECT version, snapshot_data ' +
            'FROM order_snapshots ' +
            'WHERE aggregate_id = ? ' +
            'ORDER BY version DESC LIMIT 1',
            rs -> rs.next()
                ? new Snapshot(
                    rs.getInt(1), rs.getString(2))
                : null,
            aggregateId
        );
    }

    public void saveIfNeeded(UUID aggregateId,
            OrderAggregate aggregate) {
        if (aggregate.getVersion() % SNAPSHOT_INTERVAL == 0) {
            jdbcTemplate.update(
                'INSERT INTO order_snapshots ' +
                '(aggregate_id, version, snapshot_data) ' +
                'VALUES (?, ?, ?::jsonb) ' +
                'ON CONFLICT DO NOTHING',
                aggregateId,
                aggregate.getVersion(),
                toJson(aggregate)
            );
        }
    }
}

快照间隔的选择影响两个维度:间隔太短增加存储开销和快照写入延迟,间隔太长增加回放计算量。50-100个事件是经验中的甜区。

事件溯源的适用边界与反模式

事件溯源不是银弹。以下场景应避免使用:纯查询型服务(无状态变更审计需求)、高频更新场景(如实时计数器,每次+1都产生一个事件,事件膨胀严重)、多聚合根强一致性事务(跨聚合根的事件一致性实现极其复杂)。适合使用的场景是:金融交易、订单流程、工作流引擎——任何需要“发生了什么”而非“当前是什么”的业务领域。

原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/shi-jian-su-yuan-yu-cqrs-jia-gou-shi-zhan-ding-dan-xi-tong/

(0)
小编小编
上一篇 10小时前
下一篇 10小时前

相关推荐

事件溯源与CQRS架构实战:订单系统从CRUD到事件驱动的设计演进

订单系统为何需要从CRUD演进到事件溯源

事件溯源(Event Sourcing)与CQRS(Command Query Responsibility Segregation)常被一起讨论,但它们解决的是不同维度的问题。事件溯源解决的是状态可审计性——不存当前状态,而是存导致状态变化的所有事件序列,任何时刻的状态都可从事件流重放得到。CQRS解决的是读写模型分离——写模型优化事务一致性,读模型优化查询性能。订单系统天然需要两者:订单状态变更需要完整审计轨迹(合规要求),读写负载特征差异大(写操作重一致性,读操作重聚合查询)。

事件溯源核心模型设计

传统CRUD模式下,订单状态直接写入orders表。事件溯源模式下,orders表不存在,取而代之的是order_events事件流:

-- 事件存储表
CREATE TABLE order_events (
    event_id UUID PRIMARY KEY,
    aggregate_id UUID NOT NULL,
    event_type VARCHAR(100) NOT NULL,
    event_data JSONB NOT NULL,
    version INT NOT NULL,
    created_at TIMESTAMP DEFAULT NOW(),
    UNIQUE(aggregate_id, version)
);

CREATE INDEX idx_order_events_aggregate
  ON order_events(aggregate_id, version);

订单生命周期对应的事件类型:

public class OrderEvent {
    private String eventType;
    private String aggregateId;
    private int version;
    private Map<String, Object> data;
    private LocalDateTime createdAt;
}

// 事件类型定义
// OrderCreated       - 订单创建
// ItemAdded          - 添加商品
// ItemRemoved        - 移除商品
// OrderConfirmed     - 确认订单
// PaymentReceived    - 收到付款
// OrderShipped       - 订单发货
// OrderDelivered     - 订单送达
// OrderCancelled     - 订单取消

聚合根(Aggregate)是事件溯源的核心概念,封装业务规则和状态转换逻辑:

public class OrderAggregate {
    private UUID orderId;
    private List<OrderItem> items;
    private OrderStatus status;
    private int version;

    // 从事件流重建状态
    public static OrderAggregate fromEvents(
            List<OrderEvent> events) {
        OrderAggregate order = new OrderAggregate();
        for (OrderEvent event : events) {
            order.apply(event);
        }
        return order;
    }

    private void apply(OrderEvent event) {
        switch (event.getEventType()) {
            case 'OrderCreated':
                this.orderId = UUID.fromString(
                    event.getData().get('orderId').toString());
                this.status = OrderStatus.CREATED;
                break;
            case 'ItemAdded':
                this.items.add(new OrderItem(
                    event.getData().get('sku').toString(),
                    Integer.parseInt(
                        event.getData().get('quantity').toString())
                ));
                break;
            case 'OrderConfirmed':
                this.status = OrderStatus.CONFIRMED;
                break;
            case 'PaymentReceived':
                this.status = OrderStatus.PAID;
                break;
        }
        this.version = event.getVersion();
    }

    // 业务命令 - 校验 + 产生事件
    public OrderEvent confirm() {
        if (this.status != OrderStatus.CREATED) {
            throw new IllegalStateException(
                '只有CREATED状态可确认');
        }
        if (this.items.isEmpty()) {
            throw new IllegalStateException(
                '空订单不可确认');
        }
        return OrderEvent.create(
            'OrderConfirmed', this.orderId,
            this.version + 1, Map.of());
    }
}

CQRS读写分离模型实现

写模型(Command端)直接操作事件流,读模型(Query端)通过事件投影(Projection)构建物化视图。两者数据源完全独立,通过事件总线同步:

// 写模型 - 事件存储
@Service
public class OrderCommandService {

    @Transactional
    public UUID createOrder(CreateOrderCommand cmd) {
        UUID orderId = UUID.randomUUID();
        OrderEvent event = OrderEvent.create(
            'OrderCreated', orderId, 1,
            Map.of('customerId',
                cmd.getCustomerId().toString())
        );
        eventRepository.append(event);
        eventBus.publish(event);
        return orderId;
    }
}
-- 读模型 - 物化视图
CREATE TABLE order_views (
    order_id UUID PRIMARY KEY,
    customer_id UUID,
    total_amount DECIMAL(12,2),
    status VARCHAR(20),
    item_count INT,
    updated_at TIMESTAMP
);
// 事件投影处理器 - 订阅事件更新读模型
@Component
public class OrderViewProjector {

    @EventHandler
    public void handle(OrderCreated event) {
        jdbcTemplate.update(
            'INSERT INTO order_views ' +
            '(order_id, customer_id, status, ' +
            'item_count, total_amount) ' +
            'VALUES (?, ?, \'CREATED\', 0, 0)',
            event.getAggregateId(),
            event.getData().get('customerId')
        );
    }

    @EventHandler
    public void handle(ItemAdded event) {
        jdbcTemplate.update(
            'UPDATE order_views SET ' +
            'item_count = item_count + 1, ' +
            'total_amount = total_amount + ? ' +
            'WHERE order_id = ?',
            event.getData().get('price'),
            event.getAggregateId()
        );
    }
}

快照机制优化长事件流性能

事件溯源的固有代价是状态重建开销。一个存在两年的订单如果经历了数百次状态变更,每次读取都要回放全部事件,性能不可接受。快照机制是标准解决方案:每隔N个事件保存一次聚合根的完整状态快照,重建时从最近快照开始回放。

CREATE TABLE order_snapshots (
    aggregate_id UUID,
    version INT,
    snapshot_data JSONB NOT NULL,
    created_at TIMESTAMP DEFAULT NOW(),
    PRIMARY KEY (aggregate_id, version)
);
public class SnapshotRepository {
    private static final int SNAPSHOT_INTERVAL = 50;

    public Optional<Snapshot> loadLatest(
            UUID aggregateId) {
        return jdbcTemplate.query(
            'SELECT version, snapshot_data ' +
            'FROM order_snapshots ' +
            'WHERE aggregate_id = ? ' +
            'ORDER BY version DESC LIMIT 1',
            rs -> rs.next()
                ? new Snapshot(
                    rs.getInt(1), rs.getString(2))
                : null,
            aggregateId
        );
    }

    public void saveIfNeeded(UUID aggregateId,
            OrderAggregate aggregate) {
        if (aggregate.getVersion() % SNAPSHOT_INTERVAL == 0) {
            jdbcTemplate.update(
                'INSERT INTO order_snapshots ' +
                '(aggregate_id, version, snapshot_data) ' +
                'VALUES (?, ?, ?::jsonb) ' +
                'ON CONFLICT DO NOTHING',
                aggregateId,
                aggregate.getVersion(),
                toJson(aggregate)
            );
        }
    }
}

快照间隔的选择影响两个维度:间隔太短增加存储开销和快照写入延迟,间隔太长增加回放计算量。50-100个事件是经验中的甜区。

事件溯源的适用边界与反模式

事件溯源不是银弹。以下场景应避免使用:纯查询型服务(无状态变更审计需求)、高频更新场景(如实时计数器,每次+1都产生一个事件,事件膨胀严重)、多聚合根强一致性事务(跨聚合根的事件一致性实现极其复杂)。适合使用的场景是:金融交易、订单流程、工作流引擎——任何需要“发生了什么”而非“当前是什么”的业务领域。

原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/shi-jian-su-yuan-yu-cqrs-jia-gou-shi-zhan-ding-dan-xi-tong/

(0)
小编小编
上一篇 11小时前
下一篇 11小时前

相关推荐