本文目录导读:

我将为您详细介绍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 在需要完整审计功能、复杂业务逻辑或需要时间回溯能力的系统中是非常强大的架构选择。