Event Sourcing案例

wen java案例 1

本文目录导读:

Event Sourcing案例

  1. 什么是Event Sourcing?
  2. 典型案例:电商订单系统
  3. 更复杂的案例:银行账户系统
  4. 案例:库存管理系统
  5. 实战案例:CQRS模式的订单管理
  6. 案例带来的实际价值
  7. 最佳实践建议

我将为您详细介绍Event Sourcing(事件溯源)的典型案例,从基础概念到实际应用场景。

什么是Event Sourcing?

核心思想:不直接存储对象的当前状态,而是存储导致状态变化的所有事件,当前状态可以通过重放这些事件来计算得出。

传统方式:User { id, name, email, balance }
事件溯源:UserCreated(id, name, email), BalanceChanged(id, +100), EmailChanged(id, new_email)

典型案例:电商订单系统

场景描述

用户下单、支付、发货、退货等操作,每次操作都是状态变化。

传统实现 vs 事件溯源实现

// 传统方式 - 直接存储当前状态
public class Order {
    private Long id;
    private String status;  // PENDING, PAID, SHIPPED, COMPLETED
    private BigDecimal totalAmount;
    private LocalDateTime updatedAt;
}
// 事件溯源 - 存储一系列事件
public interface OrderEvent {
    LocalDateTime getTimestamp();
}
public class OrderCreatedEvent implements OrderEvent {
    private Long orderId;
    private Long userId;
    private List<OrderItem> items;
    private BigDecimal totalAmount;
    private LocalDateTime timestamp;
}
public class OrderPaidEvent implements OrderEvent {
    private Long orderId;
    private String paymentId;
    private BigDecimal paidAmount;
    private LocalDateTime timestamp;
}
public class OrderShippedEvent implements OrderEvent {
    private Long orderId;
    private String trackingNumber;
    private LocalDateTime timestamp;
}

事件流示例

Order-123 的事件流:
1. OrderCreatedEvent (2024-01-01 10:00:00) - 创建订单
2. OrderPaidEvent    (2024-01-01 10:05:00) - 支付成功
3. OrderShippedEvent (2024-01-02 09:00:00) - 发货
4. OrderDeliveredEvent (2024-01-05 14:30:00) - 签收

更复杂的案例:银行账户系统

为什么传统方式有问题?

// 传统方式的并发问题
public class BankAccount {
    private BigDecimal balance;
    public synchronized void withdraw(BigDecimal amount) {
        if (balance.compareTo(amount) >= 0) {
            balance = balance.subtract(amount);
        }
    }
}

这种方式在并发场景下容易产生竞态条件,且无法追溯。

事件溯源实现

// 事件定义
public abstract class AccountEvent {
    private String accountId;
    private LocalDateTime timestamp;
}
public class AccountOpenedEvent extends AccountEvent {
    private String accountNumber;
    private String ownerName;
}
public class MoneyDepositedEvent extends AccountEvent {
    private BigDecimal amount;
    private String description;
}
public class MoneyWithdrawnEvent extends AccountEvent {
    private BigDecimal amount;
    private String purpose;
}
public class TransferInitiatedEvent extends AccountEvent {
    private String toAccountId;
    private BigDecimal amount;
}
// 聚合根
public class BankAccountAggregate {
    private String accountId;
    private BigDecimal balance = BigDecimal.ZERO;
    private boolean isActive = false;
    private List<AccountEvent> uncommittedEvents = new ArrayList<>();
    public static BankAccountAggregate openAccount(
            String accountId, 
            String accountNumber, 
            String ownerName) {
        BankAccountAggregate account = new BankAccountAggregate();
        AccountOpenedEvent event = new AccountOpenedEvent(accountId, accountNumber, ownerName);
        account.apply(event);
        account.uncommittedEvents.add(event);
        return account;
    }
    public void deposit(BigDecimal amount, String description) {
        if (!isActive) throw new IllegalStateException("Account not active");
        MoneyDepositedEvent event = new MoneyDepositedEvent(accountId, amount, description);
        apply(event);
        uncommittedEvents.add(event);
    }
    private void apply(AccountOpenedEvent event) {
        this.isActive = true;
        this.accountId = event.getAccountId();
    }
    private void apply(MoneyDepositedEvent event) {
        this.balance = this.balance.add(event.getAmount());
    }
    // 重放所有事件来计算当前状态
    public static BankAccountAggregate replay(List<AccountEvent> events) {
        BankAccountAggregate account = new BankAccountAggregate();
        events.forEach(account::apply);
        return account;
    }
}

案例:库存管理系统

事件定义

public class InventoryEvents {
    // 商品入库
    public record StockAdded(
        String productId, 
        int quantity, 
        String warehouseId,
        LocalDateTime timestamp
    ) {}
    // 商品出库
    public record StockDeducted(
        String productId, 
        int quantity, 
        String warehouseId,
        String orderId
    ) {}
    // 库存调整
    public record StockAdjusted(
        String productId,
        int oldQuantity,
        int newQuantity,
        String reason
    ) {}
}

重建状态的查询

public class InventoryQueryModel {
    public int getCurrentStock(String productId, List<InventoryEvent> events) {
        return events.stream()
            .filter(e -> e.getProductId().equals(productId))
            .mapToInt(event -> {
                if (event instanceof StockAdded added) {
                    return added.quantity();
                } else if (event instanceof StockDeducted deducted) {
                    return -deducted.quantity();
                }
                return 0;
            })
            .sum();
    }
    // 特定时间点的库存状态
    public int getStockAtTime(String productId, 
                              List<InventoryEvent> events, 
                              LocalDateTime pointInTime) {
        return events.stream()
            .filter(e -> e.getProductId().equals(productId))
            .filter(e -> e.timestamp().isBefore(pointInTime))
            .mapToInt(event -> {
                if (event instanceof StockAdded added) {
                    return added.quantity();
                } else if (event instanceof StockDeducted deducted) {
                    return -deducted.quantity();
                }
                return 0;
            })
            .sum();
    }
}

实战案例:CQRS模式的订单管理

完整的事件驱动架构

// 1. 命令处理端(写模型)
public class OrderCommandHandler {
    private final EventStore eventStore;
    private final EventPublisher eventPublisher;
    public void handle(CreateOrderCommand command) {
        OrderCreatedEvent event = OrderCreatedEvent.builder()
            .orderId(command.orderId())
            .userId(command.userId())
            .items(command.items())
            .status(OrderStatus.CREATED)
            .timestamp(LocalDateTime.now())
            .build();
        eventStore.saveEvents(command.orderId(), List.of(event));
        eventPublisher.publish(event);
    }
    public void handle(PayOrderCommand command) {
        // 读取当前状态
        List<DomainEvent> events = eventStore.getEventsForAggregate(command.orderId());
        OrderAggregate order = OrderAggregate.replay(events);
        // 业务校验
        if (order.getStatus() != OrderStatus.CREATED) {
            throw new IllegalStateException("订单状态不正确");
        }
        // 创建新事件
        OrderPaidEvent event = new OrderPaidEvent(command.orderId(), command.paymentInfo());
        eventStore.saveEvents(command.orderId(), List.of(event));
        eventPublisher.publish(event);
    }
}
// 2. 查询端(读模型)
@Service
public class OrderQueryService {
    private final EventStore eventStore;
    public OrderView getOrderView(String orderId) {
        List<DomainEvent> events = eventStore.getEventsForAggregate(orderId);
        OrderAggregate order = OrderAggregate.replay(events);
        return OrderView.builder()
            .orderId(order.getOrderId())
            .status(order.getStatus())
            .totalAmount(order.getTotalAmount())
            .items(order.getItems())
            .eventCount(events.size())
            .lastUpdated(events.get(events.size() - 1).getTimestamp())
            .build();
    }
    // 订单审计日志
    public List<OrderAuditEntry> getAuditTrail(String orderId) {
        return eventStore.getEventsForAggregate(orderId)
            .stream()
            .map(event -> new OrderAuditEntry(
                event.getTimestamp(),
                event.getClass().getSimpleName(),
                event.toString()
            ))
            .collect(Collectors.toList());
    }
}
// 3. 事件存储
public interface EventStore {
    void saveEvents(String aggregateId, List<DomainEvent> events);
    List<DomainEvent> getEventsForAggregate(String aggregateId);
}

案例带来的实际价值

审计与合规

// 银行交易审计
public class TransactionAuditService {
    public List<TransactionAudit> getFullTransactionHistory(String accountId) {
        return eventStore.getEventsForAggregate(accountId)
            .stream()
            .map(this::toAuditEntry)
            .collect(Collectors.toList());
    }
    // 回溯任何时间点的账户状态
    public AccountSnapshot getAccountSnapshot(String accountId, LocalDateTime time) {
        return eventStore.getEventsForAggregate(accountId)
            .stream()
            .filter(e -> e.getTimestamp().isBefore(time))
            .reduce(AccountSnapshot.empty(), this::applyEventToSnapshot, (s1, s2) -> s2);
    }
}

系统故障恢复

public class RecoveryService {
    private final EventStore eventStore;
    public void recoverAccountState() {
        // 获取所有未处理的事件
        List<DomainEvent> unrecoveredEvents = 
            eventStore.getEventsAfterLastLogCheckpoint();
        // 重放事件恢复状态
        AccountAggregate account = AccountAggregate.replay(unrecoveredEvents);
        // 更新读模型
        updateReadModel(account);
        updateSearchIndex(account);
        updateAnalytics(account);
    }
}

最佳实践建议

事件设计准则

// 错误设计 - 事件包含太多状态
public record OrderEvent(
    Long orderId,
    String orderStatus,  // 不要存完整状态
    BigDecimal totalAmount, // 不要存计算结果
    List<OrderItem> items, // 不要存整个聚合
    String customerName,
    String customerEmail
) {}
// 正确设计 - 事件只包含变化的信息
public record OrderStatusChanged(
    Long orderId,
    OrderStatus from,
    OrderStatus to,
    String changedBy,
    String reason
) {}

性能优化策略

public class OptimizedEventStore {
    // 定期快照减少重放时间
    public Snapshot createSnapshot(String aggregateId) {
        List<DomainEvent> events = getEventsAfterLastSnapshot(aggregateId);
        AggregateState state = AggregateState.replay(events);
        return new Snapshot(aggregateId, state, 
            getLastEventVersion(aggregateId));
    }
    // 从快照 + 增量事件恢复
    public AggregateState restoreAggregate(String aggregateId) {
        Snapshot snapshot = getLatestSnapshot(aggregateId);
        List<DomainEvent> events = getEventsAfterVersion(
            aggregateId, snapshot.getVersion());
        return AggregateState.replay(snapshot.getState(), events);
    }
}

Event Sourcing 的核心价值:

  • 完整的审计轨迹 - 每个操作都有记录
  • 时间旅行 - 可以回到任何时间点
  • 可重放性 - 支持故障恢复和调试
  • 解耦系统组件 - 事件作为系统间的通信方式

适合场景:

  • 金融、银行系统
  • 电商平台
  • 供应链管理
  • 审计和合规要求高的系统

需要权衡的方面:

  • 查询性能可能较低(需要事件重放)
  • 存储需求更大(存储事件流)
  • 系统复杂度增加
  • 需要处理事件版本兼容

Event Sourcing 在需要完整审计功能、复杂业务逻辑或需要时间回溯能力的系统中是非常强大的架构选择。

抱歉,评论功能暂时关闭!