一、架构演进:为什么需要事件驱动

传统后端架构中,系统状态通常以CRUD方式直接存储在关系数据库中。这种方式在单体应用时代运转良好,但在复杂业务系统和高并发场景下逐渐暴露出根本性缺陷:数据覆盖导致历史丢失、读写争抢导致性能瓶颈、紧耦合导致扩展困难

事件驱动架构(Event-Driven Architecture)及其核心模式——事件溯源(Event Sourcing)命令查询职责分离(CQRS)——正是为了解决这些问题而生。它们将"状态变化"提升为一等公民,而非仅仅保存当前状态的快照。

Martin Fowler曾指出:"事件溯源是一种看似简单却影响深远的架构决策——它改变了我们对数据持久化的认知方式。"在金融交易、电商订单、物联网、协同编辑等需要完整审计追踪和时间旅行能力的场景中,这一模式正得到越来越广泛的应用。

二、核心概念解析

2.1 命令查询职责分离(CQRS)

CQRS由Greg Young在2010年提出,其核心思想是将对数据的修改操作(Command)和查询操作(Query)使用不同的模型来处理

  • Command Model(写模型):关注业务规则验证、领域逻辑执行和一致性保证,负责产生领域事件
  • Query Model(读模型):关注数据展示优化、查询性能提升,可以针对不同查询场景独立设计数据结构

传统架构中读写共享同一数据模型,高并发下必然导致锁竞争和性能瓶颈。CQRS通过物理或逻辑分离,让读模型可以为查询场景做极致优化(如宽表、物化视图、搜索引擎索引),写模型则可以保持面向领域的纯净结构。

2.2 事件溯源(Event Sourcing)

事件溯源颠覆了传统CRUD的持久化方式:不直接存储对象的当前状态,而是存储导致状态变化的完整事件序列

以一个账户系统为例:

  • 传统方式:balance字段直接从100更新为80,历史信息丢失
  • 事件溯源:存储AccountCreated(初始100)MoneyWithdrawn(取出20)两个事件,当前状态通过事件回放推导

这带来三个核心能力:

  1. 完整审计追踪:每一次状态变化都有据可查,天然满足合规要求
  2. 时间旅行查询:可重建任意历史时刻的系统状态
  3. 事件重放调试:可在测试环境重放生产事件流复现问题

2.3 CQRS + Event Sourcing = 黄金组合

两者结合形成天然互补:Event Sourcing为CQRS提供理想的数据源——写模型产生事件并持久化到事件存储(Event Store),事件通过异步管道投影到多个读模型。读模型可以针对不同查询场景自由设计,从SQL宽表到Elasticsearch全文索引均可。

三、架构设计与核心组件

3.1 整体架构图

一个典型的CQRS + Event Sourcing系统包含以下核心组件:客户端发送命令到命令总线,命令处理器加载聚合根历史事件、执行业务逻辑、产生新事件写入Event Store,事件总线将新事件异步分发给各投影器(Projector),投影器将事件投影到读模型存储(如PostgreSQL、Redis、Elasticsearch),查询处理器直接从读模型提供数据。

3.2 事件设计原则

事件是系统中唯一的事实来源,其设计质量直接决定系统的可维护性和演进能力:

  • 不可变性:事件一旦产生就不可修改,只能产生新事件来修正
  • 自描述性:事件名称使用过去式动词(如OrderPlaced、PaymentConfirmed),包含足够的业务语义
  • 幂等消费:事件可能被重复投递,消费者必须能正确处理重复事件
  • 版本兼容:事件Schema演化需要考虑向后兼容(增添字段、向上转换)

3.3 聚合根设计

聚合根( Aggregate Root )是CQRS + ES架构中写模型的核心,负责维护业务不变量并产生事件:

  • 唯一标识聚合实例,是事件流的基本单位
  • 接收命令后执行业务规则校验,通过后产生新事件
  • 通过应用事件来更新内部状态(状态是事件的函数)
  • 每次持久化时进行乐观并发检查(版本号机制保证线性一致性)

四、Java + Spring Boot 实战实现

4.1 核心领域模型

// 领域事件基类
public abstract class DomainEvent {
    private final String eventId = UUID.randomUUID().toString();
    private final String aggregateId;
    private final long timestamp = System.currentTimeMillis();
    private final int version;
    
    protected DomainEvent(String aggregateId, int version) {
        this.aggregateId = aggregateId;
        this.version = version;
    }
    
    // getters...
}

// 订单相关事件
@Data
@AllArgsConstructor
public class OrderCreatedEvent extends DomainEvent {
    private String orderId;
    private String customerId;
    private List items;
    private BigDecimal totalAmount;
    
    public OrderCreatedEvent(String aggregateId, int version, 
            String customerId, List items, BigDecimal totalAmount) {
        super(aggregateId, version);
        this.orderId = aggregateId;
        this.customerId = customerId;
        this.items = items;
        this.totalAmount = totalAmount;
    }
}

@Data
@AllArgsConstructor
public class OrderPaidEvent extends DomainEvent {
    private String orderId;
    private String paymentId;
    private BigDecimal paidAmount;
    private LocalDateTime paidAt;
    
    public OrderPaidEvent(String aggregateId, int version,
            String paymentId, BigDecimal paidAmount, LocalDateTime paidAt) {
        super(aggregateId, version);
        this.orderId = aggregateId;
        this.paymentId = paymentId;
        this.paidAmount = paidAmount;
        this.paidAt = paidAt;
    }
}

@Data
@AllArgsConstructor
public class OrderShippedEvent extends DomainEvent {
    private String orderId;
    private String trackingNumber;
    private String carrier;
    
    public OrderShippedEvent(String aggregateId, int version,
            String trackingNumber, String carrier) {
        super(aggregateId, version);
        this.orderId = aggregateId;
        this.trackingNumber = trackingNumber;
        this.carrier = carrier;
    }
}

@Data
@AllArgsConstructor
public class OrderCancelledEvent extends DomainEvent {
    private String orderId;
    private String reason;
    
    public OrderCancelledEvent(String aggregateId, int version, String reason) {
        super(aggregateId, version);
        this.orderId = aggregateId;
        this.reason = reason;
    }
}

// 值对象
@Data
@AllArgsConstructor
public class OrderItem {
    private String productId;
    private String productName;
    private int quantity;
    private BigDecimal unitPrice;
}

// 订单状态枚举
public enum OrderStatus {
    CREATED, PAID, SHIPPED, CANCELLED
}

4.2 聚合根实现

@Getter
public class OrderAggregate {
    
    private String orderId;
    private String customerId;
    private List items;
    private BigDecimal totalAmount;
    private OrderStatus status;
    private String paymentId;
    private String trackingNumber;
    private String carrier;
    private String cancelReason;
    private int version;
    
    // 待提交的新事件列表
    @Getter(AccessLevel.NONE)
    private final List uncommittedEvents = new ArrayList<>();
    
    // 私有构造函数,通过工厂方法创建
    private OrderAggregate() {}
    
    /**
     * 工厂方法:创建新订单
     */
    public static OrderAggregate create(String customerId, List items) {
        OrderAggregate aggregate = new OrderAggregate();
        
        // 业务规则校验
        if (items == null || items.isEmpty()) {
            throw new IllegalArgumentException("订单必须包含至少一个商品");
        }
        
        BigDecimal total = items.stream()
            .map(item -> item.getUnitPrice().multiply(BigDecimal.valueOf(item.getQuantity())))
            .reduce(BigDecimal.ZERO, BigDecimal::add);
        
        // 产生领域事件
        String orderId = UUID.randomUUID().toString();
        aggregate.applyEvent(new OrderCreatedEvent(orderId, 0, customerId, items, total));
        
        return aggregate;
    }
    
    /**
     * 支付订单
     */
    public void pay(String paymentId, BigDecimal amount) {
        if (this.status != OrderStatus.CREATED) {
            throw new IllegalStateException("只能支付已创建的订单,当前状态: " + this.status);
        }
        if (this.totalAmount.compareTo(amount) != 0) {
            throw new IllegalArgumentException("支付金额与订单金额不符");
        }
        applyEvent(new OrderPaidEvent(this.orderId, this.version, paymentId, amount, LocalDateTime.now()));
    }
    
    /**
     * 发货
     */
    public void ship(String trackingNumber, String carrier) {
        if (this.status != OrderStatus.PAID) {
            throw new IllegalStateException("只能对已付款的订单发货");
        }
        if (trackingNumber == null || trackingNumber.isBlank()) {
            throw new IllegalArgumentException("物流单号不能为空");
        }
        applyEvent(new OrderShippedEvent(this.orderId, this.version, trackingNumber, carrier));
    }
    
    /**
     * 取消订单
     */
    public void cancel(String reason) {
        if (this.status == OrderStatus.SHIPPED) {
            throw new IllegalStateException("已发货的订单不能直接取消");
        }
        if (this.status == OrderStatus.CANCELLED) {
            throw new IllegalStateException("订单已取消,不能重复取消");
        }
        applyEvent(new OrderCancelledEvent(this.orderId, this.version, reason));
    }
    
    /**
     * 事件溯源:从历史事件重建聚合状态
     */
    public static OrderAggregate rehydrate(String orderId, List history) {
        OrderAggregate aggregate = new OrderAggregate();
        for (DomainEvent event : history) {
            aggregate.applyEvent(event);
        }
        return aggregate;
    }
    
    /**
     * 应用事件 - 更新内部状态
     */
    private void applyEvent(DomainEvent event) {
        if (event instanceof OrderCreatedEvent e) {
            this.orderId = e.getOrderId();
            this.customerId = e.getCustomerId();
            this.items = e.getItems();
            this.totalAmount = e.getTotalAmount();
            this.status = OrderStatus.CREATED;
        } else if (event instanceof OrderPaidEvent e) {
            this.paymentId = e.getPaymentId();
            this.status = OrderStatus.PAID;
        } else if (event instanceof OrderShippedEvent e) {
            this.trackingNumber = e.getTrackingNumber();
            this.carrier = e.getCarrier();
            this.status = OrderStatus.SHIPPED;
        } else if (event instanceof OrderCancelledEvent e) {
            this.cancelReason = e.getReason();
            this.status = OrderStatus.CANCELLED;
        }
        
        this.version = event.getVersion() + 1;
        
        // 如果是新事件(不是重放历史事件),记录到待提交列表
        if (event.getVersion() == this.version - 1 && 
            event.getVersion() == uncommittedEvents.size()) {
            uncommittedEvents.add(event);
        }
    }
    
    /**
     * 获取未提交的事件并清空列表
     */
    public List getUncommittedEvents() {
        List events = new ArrayList<>(uncommittedEvents);
        uncommittedEvents.clear();
        return events;
    }
}

4.3 事件存储(Event Store)

/**
 * 事件存储接口 - 定义事件持久化的核心操作
 * 生产环境可使用EventStoreDB、Axon Framework或自建实现
 */
public interface EventStore {
    
    /**
     * 追加事件到指定聚合的事件流
     * @param aggregateId 聚合根ID
     * @param events 待追加的事件列表
     * @param expectedVersion 期望的当前版本(乐观并发控制)
     */
    void appendEvents(String aggregateId, List events, int expectedVersion);
    
    /**
     * 读取指定聚合的所有事件
     */
    List readEvents(String aggregateId);
    
    /**
     * 从指定版本开始读取事件(用于快照优化后的增量加载)
     */
    List readEvents(String aggregateId, int fromVersion);
}

/**
 * 基于关系数据库的事件存储实现
 */
@Repository
@RequiredArgsConstructor
public class JdbcEventStore implements EventStore {
    
    private final JdbcTemplate jdbcTemplate;
    private final ObjectMapper objectMapper;
    
    @Override
    @Transactional
    public void appendEvents(String aggregateId, List events, int expectedVersion) {
        // 乐观并发检查
        Integer currentVersion = jdbcTemplate.queryForObject(
            "SELECT MAX(version) FROM event_store WHERE aggregate_id = ?",
            Integer.class, aggregateId
        );
        currentVersion = currentVersion == null ? -1 : currentVersion;
        
        if (currentVersion != expectedVersion) {
            throw new ConcurrencyException(
                "并发冲突: 聚合 " + aggregateId + " 当前版本 " + currentVersion + ",期望版本 " + expectedVersion
            );
        }
        
        // 批量写入事件
        String sql = """
            INSERT INTO event_store (aggregate_id, version, event_type, event_data, occurred_at, event_id)
            VALUES (?, ?, ?, ?::jsonb, ?, ?)
        """;
        
        int version = expectedVersion;
        for (DomainEvent event : events) {
            version++;
            try {
                String eventData = objectMapper.writeValueAsString(event);
                jdbcTemplate.update(sql,
                    aggregateId,
                    version,
                    event.getClass().getSimpleName(),
                    eventData,
                    Instant.ofEpochMilli(event.getTimestamp()),
                    event.getEventId()
                );
            } catch (JsonProcessingException e) {
                throw new RuntimeException("事件序列化失败", e);
            }
        }
    }
    
    @Override
    public List readEvents(String aggregateId) {
        return readEvents(aggregateId, 0);
    }
    
    @Override
    public List readEvents(String aggregateId, int fromVersion) {
        String sql = """
            SELECT event_type, event_data, version 
            FROM event_store 
            WHERE aggregate_id = ? AND version >= ?
            ORDER BY version ASC
        """;
        
        return jdbcTemplate.query(sql, (rs, rowNum) -> {
            String eventType = rs.getString("event_type");
            String eventData = rs.getString("event_data");
            try {
                return deserializeEvent(eventType, eventData);
            } catch (Exception e) {
                throw new RuntimeException("事件反序列化失败: " + eventType, e);
            }
        }, aggregateId, fromVersion);
    }
    
    private DomainEvent deserializeEvent(String eventType, String data) throws JsonProcessingException {
        return switch (eventType) {
            case "OrderCreatedEvent" -> objectMapper.readValue(data, OrderCreatedEvent.class);
            case "OrderPaidEvent" -> objectMapper.readValue(data, OrderPaidEvent.class);
            case "OrderShippedEvent" -> objectMapper.readValue(data, OrderShippedEvent.class);
            case "OrderCancelledEvent" -> objectMapper.readValue(data, OrderCancelledEvent.class);
            default -> throw new IllegalArgumentException("未知事件类型: " + eventType);
        };
    }
}

/**
 * 建表SQL - PostgreSQL
 */
// CREATE TABLE event_store (
//     id BIGSERIAL PRIMARY KEY,
//     aggregate_id VARCHAR(64) NOT NULL,
//     version INT NOT NULL,
//     event_type VARCHAR(128) NOT NULL,
//     event_data JSONB NOT NULL,
//     occurred_at TIMESTAMP NOT NULL DEFAULT NOW(),
//     event_id VARCHAR(64) NOT NULL UNIQUE,
//     UNIQUE (aggregate_id, version)
// );
// 
// CREATE INDEX idx_event_store_aggregate ON event_store(aggregate_id, version);
// CREATE INDEX idx_event_store_occurred ON event_store(occurred_at);

4.4 命令处理器与读模型投影

/**
 * 命令处理器 - 接收命令,协调聚合根完成业务操作
 */
@Service
@RequiredArgsConstructor
public class OrderCommandHandler {
    
    private final EventStore eventStore;
    private final EventPublisher eventPublisher;
    
    /**
     * 处理创建订单命令
     */
    public String handle(CreateOrderCommand command) {
        // 1. 创建聚合根
        OrderAggregate order = OrderAggregate.create(
            command.getCustomerId(), command.getItems()
        );
        
        // 2. 持久化事件(乐观并发控制初始版本为-1)
        List events = order.getUncommittedEvents();
        eventStore.appendEvents(order.getOrderId(), events, -1);
        
        // 3. 发布事件到消息总线(异步通知读模型和其他服务)
        eventPublisher.publish(events);
        
        return order.getOrderId();
    }
    
    /**
     * 处理支付命令
     */
    @Transactional
    public void handle(PayOrderCommand command) {
        // 1. 从事件存储重建聚合根
        List history = eventStore.readEvents(command.getOrderId());
        if (history.isEmpty()) {
            throw new OrderNotFoundException(command.getOrderId());
        }
        
        OrderAggregate order = OrderAggregate.rehydrate(command.getOrderId(), history);
        
        // 2. 执行业务操作
        order.pay(command.getPaymentId(), command.getAmount());
        
        // 3. 持久化新事件(乐观并发检查当前版本)
        List newEvents = order.getUncommittedEvents();
        int currentVersion = history.get(history.size() - 1).getVersion();
        eventStore.appendEvents(command.getOrderId(), newEvents, currentVersion);
        
        // 4. 发布事件
        eventPublisher.publish(newEvents);
    }
}

/**
 * 读模型投影器 - 将事件投影到读模型(订单列表视图)
 */
@Component
@RequiredArgsConstructor
@Slf4j
public class OrderViewProjector {
    
    private final JdbcTemplate jdbcTemplate;
    
    /**
     * 处理OrderCreated事件 - 插入读模型记录
     */
    @EventListener
    public void on(OrderCreatedEvent event) {
        String sql = """
            INSERT INTO order_view (order_id, customer_id, total_amount, status, created_at, version)
            VALUES (?, ?, ?, 'CREATED', ?, 1)
            ON CONFLICT (order_id) DO UPDATE SET
                customer_id = EXCLUDED.customer_id,
                total_amount = EXCLUDED.total_amount,
                status = 'CREATED',
                created_at = EXCLUDED.created_at,
                version = EXCLUDED.version
        """;
        
        jdbcTemplate.update(sql,
            event.getOrderId(), event.getCustomerId(),
            event.getTotalAmount(), Instant.ofEpochMilli(event.getTimestamp())
        );
        
        log.info("读模型已更新: 订单 {} 已创建", event.getOrderId());
    }
    
    /**
     * 处理OrderPaid事件 - 更新读模型状态
     */
    @EventListener
    public void on(OrderPaidEvent event) {
        String sql = """
            UPDATE order_view 
            SET status = 'PAID', payment_id = ?, paid_at = ?, version = ?
            WHERE order_id = ?
        """;
        
        jdbcTemplate.update(sql,
            event.getPaymentId(),
            event.getPaidAt(),
            event.getVersion() + 1,
            event.getOrderId()
        );
        
        log.info("读模型已更新: 订单 {} 已支付", event.getOrderId());
    }
    
    @EventListener
    public void on(OrderShippedEvent event) {
        jdbcTemplate.update("""
            UPDATE order_view 
            SET status = 'SHIPPED', tracking_number = ?, carrier = ?, shipped_at = NOW(), version = ?
            WHERE order_id = ?
        """, event.getTrackingNumber(), event.getCarrier(), 
            event.getVersion() + 1, event.getOrderId());
    }
    
    @EventListener
    public void on(OrderCancelledEvent event) {
        jdbcTemplate.update("""
            UPDATE order_view 
            SET status = 'CANCELLED', cancel_reason = ?, cancelled_at = NOW(), version = ?
            WHERE order_id = ?
        """, event.getReason(), event.getVersion() + 1, event.getOrderId());
    }
}

/**
 * 读模型查询服务
 */
@Service
@RequiredArgsConstructor
public class OrderQueryService {
    
    private final JdbcTemplate jdbcTemplate;
    
    /**
     * 查询订单详情 - 从读模型直接读取(高性能)
     */
    public OrderViewDTO getOrder(String orderId) {
        return jdbcTemplate.queryForObject("""
            SELECT order_id, customer_id, total_amount, status, 
                   payment_id, tracking_number, carrier, 
                   created_at, paid_at, shipped_at, cancelled_at
            FROM order_view WHERE order_id = ?
        """, (rs, rowNum) -> OrderViewDTO.builder()
            .orderId(rs.getString("order_id"))
            .customerId(rs.getString("customer_id"))
            .totalAmount(rs.getBigDecimal("total_amount"))
            .status(rs.getString("status"))
            .trackingNumber(rs.getString("tracking_number"))
            .carrier(rs.getString("carrier"))
            .createdAt(rs.getTimestamp("created_at").toInstant())
            .build()
        , orderId);
    }
    
    /**
     * 分页查询用户订单列表
     */
    public List listCustomerOrders(String customerId, int page, int size) {
        return jdbcTemplate.query("""
            SELECT order_id, customer_id, total_amount, status, tracking_number, created_at
            FROM order_view 
            WHERE customer_id = ? 
            ORDER BY created_at DESC 
            LIMIT ? OFFSET ?
        """, (rs, rowNum) -> OrderViewDTO.builder()
                .orderId(rs.getString("order_id"))
                .customerId(rs.getString("customer_id"))
                .totalAmount(rs.getBigDecimal("total_amount"))
                .status(rs.getString("status"))
                .build(),
            customerId, size, page * size
        );
    }
}

// 读模型建表SQL
// CREATE TABLE order_view (
//     order_id VARCHAR(64) PRIMARY KEY,
//     customer_id VARCHAR(64) NOT NULL,
//     total_amount DECIMAL(12,2) NOT NULL,
//     status VARCHAR(16) NOT NULL,
//     payment_id VARCHAR(64),
//     tracking_number VARCHAR(64),
//     carrier VARCHAR(32),
//     cancel_reason TEXT,
//     created_at TIMESTAMP NOT NULL,
//     paid_at TIMESTAMP,
//     shipped_at TIMESTAMP,
//     cancelled_at TIMESTAMP,
//     version INT NOT NULL DEFAULT 0
// );
// CREATE INDEX idx_order_view_customer ON order_view(customer_id, created_at DESC);

五、高级主题

5.1 快照机制(Snapshot)

事件溯源的一个核心挑战是:当事件流很长时,每次重建聚合都从第一个事件回放效率很低。快照(Snapshot)机制通过在特定版本点保存聚合状态的完整拷贝,让重建只需从最近快照开始回放其后的事件:

@Repository
@RequiredArgsConstructor
public class SnapshotEventStore implements EventStore {
    
    private final JdbcTemplate jdbcTemplate;
    private final EventStore eventStore;
    private final ObjectMapper objectMapper;
    
    private static final int SNAPSHOT_FREQUENCY = 20; // 每20个事件做一次快照
    
    @Override
    public List readEvents(String aggregateId) {
        // 1. 读取最新快照
        Snapshot snapshot = loadSnapshot(aggregateId);
        
        if (snapshot == null) {
            // 无快照,读取全部事件
            return eventStore.readEvents(aggregateId);
        }
        
        // 2. 读取快照之后的事件
        List eventsAfterSnapshot = eventStore.readEvents(
            aggregateId, snapshot.getVersion() + 1
        );
        
        return eventsAfterSnapshot;
    }
    
    @Override
    public void appendEvents(String aggregateId, List events, int expectedVersion) {
        eventStore.appendEvents(aggregateId, events, expectedVersion);
        
        int newVersion = expectedVersion + events.size();
        
        // 达到快照频率时,自动生成快照
        if (newVersion % SNAPSHOT_FREQUENCY == 0) {
            List allEvents = eventStore.readEvents(aggregateId);
            OrderAggregate aggregate = OrderAggregate.rehydrate(aggregateId, allEvents);
            saveSnapshot(aggregateId, newVersion, aggregate);
        }
    }
    
    private void saveSnapshot(String aggregateId, int version, OrderAggregate aggregate) {
        try {
            String stateJson = objectMapper.writeValueAsString(aggregate);
            jdbcTemplate.update("""
                INSERT INTO snapshots (aggregate_id, version, state_data, created_at)
                VALUES (?, ?, ?::jsonb, NOW())
                ON CONFLICT (aggregate_id) DO UPDATE SET
                    version = EXCLUDED.version,
                    state_data = EXCLUDED.state_data,
                    created_at = EXCLUDED.created_at
            """, aggregateId, version, stateJson);
        } catch (Exception e) {
            // 快照保存失败不影响主流程
        }
    }
    
    private Snapshot loadSnapshot(String aggregateId) {
        try {
            return jdbcTemplate.queryForObject(
                "SELECT version, state_data FROM snapshots WHERE aggregate_id = ?",
                (rs, rowNum) -> new Snapshot(
                    rs.getInt("version"), rs.getString("state_data")
                ),
                aggregateId
            );
        } catch (EmptyResultDataAccessException e) {
            return null;
        }
    }
}

5.2 事件版本兼容与Schema演化

随着业务发展,事件Schema可能需要变更。常见策略包括:

  • 向上转换(Upcasting):读取旧事件时通过转换器转换为新格式,v1.OrderCreatedEvent向上转换为v2.OrderCreatedEvent(如新增字段则赋予默认值)
  • 双写过渡:新版本应用写入时同时写两种格式事件,迁移完成后废弃旧格式
  • 版本标记:事件携带version字段,反序列化时根据版本号选择不同的构造逻辑

5.3 CAP定理下的权衡

CQRS + Event Sourcing本质上是一种最终一致性架构。写模型产生的事件通过异步管道同步到读模型,存在毫秒到秒级的延迟。这意味着:

  • 写后读一致性:用户操作后立即查询可能看不到最新结果,需要额外策略(如命令返回版本号、客户端轮询、WebSocket推送)
  • 适合场景:对实时性要求不高但需要完整审计能力的场景(订单管理、金融账务、库存管理)
  • 不适合场景:需要强一致性的实时库存扣减、秒杀库存校验等可在写模型侧用聚合根直接校验

5.4 与消息队列的集成

事件发布通常集成Kafka、RabbitMQ等消息中间件:

@Component
@RequiredArgsConstructor
public class KafkaEventPublisher implements EventPublisher {
    
    private final KafkaTemplate kafkaTemplate;
    private final ObjectMapper objectMapper;
    
    @Override
    public void publish(List events) {
        for (DomainEvent event : events) {
            String topic = "domain-events." + event.getClass().getSimpleName();
            String key = event.getAggregateId();
            
            try {
                String payload = objectMapper.writeValueAsString(event);
                kafkaTemplate.send(topic, key, payload)
                    .whenComplete((result, ex) -> {
                        if (ex != null) {
                            log.error("事件发布失败: {}", event.getEventId(), ex);
                        }
                    });
            } catch (JsonProcessingException e) {
                log.error("事件序列化失败", e);
            }
        }
    }
}

消费者侧需注意幂等处理——Kafka的at-least-once语义意味着事件可能被重复消费,可通过事件ID去重或设计天然幂等的投影逻辑来解决。

六、方案对比与选型决策

CQRS + Event Sourcing并非万能方案,需要根据实际场景权衡。与常见架构模式对比:传统CRUD实现简单但无法追溯历史;事件溯源与CQRS提供完整审计能力和强大扩展性但实现复杂度最高;事件日志加CRUD是折中方案。

七、生产部署建议

  1. 从非核心系统开始:先用在新模块或B端管理工具上积累经验,再推广到核心交易链路
  2. Event Store选型:小规模可自建PostgreSQL事件表;大规模建议EventStoreDB、Axon Server或Kafka作为事件日志
  3. 监控与告警:重点关注事件写入延迟、投影延迟、死信队列堆积、版本冲突频率
  4. 事件不可变是底线:绝对不允许修改已发布的事件,纠错只能产生补偿事件
  5. 幂等骨架:所有消费者从一开始就设计为幂等的,这对系统长期可维护性至关重要
  6. 渐进式演进:不要试图一步到位ES化整个系统,可从关键聚合开始,周边仍用CRUD
  7. 团队认知对齐:团队成员需要理解最终一致性和领域事件思维,这需要培训和沟通成本

八、总结

事件溯源与CQRS是一对天然组合,它们将"状态变化"提升为一等公民,赋予系统完整审计、时间旅行和灵活投影的能力。虽然带来了一定的实现复杂性和最终一致性的挑战,但在需要深度业务洞察和无限扩展可能的复杂业务系统中,这种架构模式正成为越来越多团队的首选。

核心理念可以概括为:事件是唯一的事实来源,状态是事件的投影,而投影可以无限多样。掌握这一思想,就能构建出既保持数据完整性的还原能力、又拥有极致查询弹性的现代化后端系统。

点赞(0) 打赏

评论列表 共有 0 条评论

暂无评论
立即
投稿

微信公众账号

微信扫一扫加关注

发表
评论
返回
顶部