本文目录导读:

我将为您提供一个完整的Java事件溯源(Event Sourcing)案例,以银行账户系统为例:
事件定义
// 基础事件接口
public interface DomainEvent {
String getEventId();
LocalDateTime getOccurredAt();
int getVersion();
}
// 账户事件基类
public abstract class AccountEvent implements DomainEvent {
private final String eventId;
private final LocalDateTime occurredAt;
private final int version;
public AccountEvent(String eventId, LocalDateTime occurredAt, int version) {
this.eventId = eventId;
this.occurredAt = occurredAt;
this.version = version;
}
@Override
public String getEventId() { return eventId; }
@Override
public LocalDateTime getOccurredAt() { return occurredAt; }
@Override
public int getVersion() { return version; }
}
// 账户创建事件
public class AccountCreatedEvent extends AccountEvent {
private final String accountId;
private final String accountHolderName;
private final BigDecimal initialBalance;
public AccountCreatedEvent(String eventId, String accountId,
String accountHolderName, BigDecimal initialBalance,
LocalDateTime occurredAt, int version) {
super(eventId, occurredAt, version);
this.accountId = accountId;
this.accountHolderName = accountHolderName;
this.initialBalance = initialBalance;
}
// Getters...
}
// 存款事件
public class MoneyDepositedEvent extends AccountEvent {
private final String accountId;
private final BigDecimal amount;
// 构造方法、getters...
}
// 取款事件
public class MoneyWithdrawnEvent extends AccountEvent {
private final String accountId;
private final BigDecimal amount;
// 构造方法、getters...
}
// 账户冻结事件
public class AccountFrozenEvent extends AccountEvent {
private final String accountId;
private final String reason;
// 构造方法、getters...
}
聚合根
public class BankAccount {
private String accountId;
private String accountHolderName;
private BigDecimal balance;
private boolean isFrozen;
private int version;
private List<DomainEvent> pendingEvents = new ArrayList<>();
// 构造函数 - 创建新账户
public static BankAccount createAccount(String accountId, String holderName) {
if (accountId == null || accountId.isEmpty()) {
throw new IllegalArgumentException("Account ID cannot be null or empty");
}
if (holderName == null || holderName.isEmpty()) {
throw new IllegalArgumentException("Account holder name cannot be null or empty");
}
BankAccount account = new BankAccount();
AccountCreatedEvent event = new AccountCreatedEvent(
UUID.randomUUID().toString(),
accountId,
holderName,
BigDecimal.ZERO.setScale(2),
LocalDateTime.now(),
1
);
account.apply(event);
return account;
}
// 恢复聚合(从事件重建)
public static BankAccount rebuild(List<AccountEvent> events) {
BankAccount account = new BankAccount();
events.forEach(account::apply);
return account;
}
// 业务方法 - 存款
public void deposit(BigDecimal amount) {
validateOperation();
if (amount == null || amount.compareTo(BigDecimal.ZERO) <= 0) {
throw new IllegalArgumentException("Deposit amount must be positive");
}
MoneyDepositedEvent event = new MoneyDepositedEvent(
UUID.randomUUID().toString(),
this.accountId,
amount,
LocalDateTime.now(),
this.version + 1
);
apply(event);
}
// 业务方法 - 取款
public void withdraw(BigDecimal amount) {
validateOperation();
if (amount == null || amount.compareTo(BigDecimal.ZERO) <= 0) {
throw new IllegalArgumentException("Withdrawal amount must be positive");
}
if (this.balance.compareTo(amount) < 0) {
throw new IllegalStateException("Insufficient funds");
}
MoneyWithdrawnEvent event = new MoneyWithdrawnEvent(
UUID.randomUUID().toString(),
this.accountId,
amount,
LocalDateTime.now(),
this.version + 1
);
apply(event);
}
// 业务方法 - 冻结账户
public void freeze(String reason) {
if (this.isFrozen) {
throw new IllegalStateException("Account is already frozen");
}
AccountFrozenEvent event = new AccountFrozenEvent(
UUID.randomUUID().toString(),
this.accountId,
reason,
LocalDateTime.now(),
this.version + 1
);
apply(event);
}
// 事件应用逻辑
private void apply(DomainEvent event) {
if (event instanceof AccountCreatedEvent) {
apply((AccountCreatedEvent) event);
} else if (event instanceof MoneyDepositedEvent) {
apply((MoneyDepositedEvent) event);
} else if (event instanceof MoneyWithdrawnEvent) {
apply((MoneyWithdrawnEvent) event);
} else if (event instanceof AccountFrozenEvent) {
apply((AccountFrozenEvent) event);
}
this.version = event.getVersion();
this.pendingEvents.add(event);
}
private void apply(AccountCreatedEvent event) {
this.accountId = event.getAccountId();
this.accountHolderName = event.getAccountHolderName();
this.balance = event.getInitialBalance();
this.isFrozen = false;
}
private void apply(MoneyDepositedEvent event) {
this.balance = this.balance.add(event.getAmount());
}
private void apply(MoneyWithdrawnEvent event) {
this.balance = this.balance.subtract(event.getAmount());
}
private void apply(AccountFrozenEvent event) {
this.isFrozen = true;
}
private void validateOperation() {
if (this.isFrozen) {
throw new IllegalStateException("Account is frozen");
}
}
// Getter方法
public List<DomainEvent> getPendingEvents() {
return List.copyOf(pendingEvents);
}
public void clearPendingEvents() {
this.pendingEvents.clear();
}
public int getVersion() { return version; }
public String getAccountId() { return accountId; }
// 其他getters...
}
事件存储
// 事件存储接口
public interface EventStore {
void append(List<DomainEvent> events);
List<DomainEvent> load(String aggregateId);
List<DomainEvent> loadSince(String aggregateId, int version);
}
// 内存事件存储实现
public class InMemoryEventStore implements EventStore {
private final Map<String, List<DomainEvent>> eventMap = new ConcurrentHashMap<>();
@Override
public synchronized void append(List<DomainEvent> events) {
if (events.isEmpty()) return;
String aggregateId = getAggregateId(events.get(0));
List<DomainEvent> existingEvents = eventMap.getOrDefault(aggregateId, new ArrayList<>());
int expectedVersion = existingEvents.size();
DomainEvent firstEvent = events.get(0);
if (firstEvent.getVersion() != expectedVersion + 1) {
throw new ConcurrencyException("Optimistic lock conflict");
}
existingEvents.addAll(events);
eventMap.put(aggregateId, existingEvents);
}
@Override
public List<DomainEvent> load(String aggregateId) {
return eventMap.getOrDefault(aggregateId, Collections.emptyList());
}
@Override
public List<DomainEvent> loadSince(String aggregateId, int version) {
List<DomainEvent> events = load(aggregateId);
return events.stream()
.filter(event -> event.getVersion() > version)
.collect(Collectors.toList());
}
private String getAggregateId(DomainEvent event) {
if (event instanceof AccountCreatedEvent) {
return ((AccountCreatedEvent) event).getAccountId();
} else if (event instanceof MoneyDepositedEvent) {
return ((MoneyDepositedEvent) event).getAccountId();
}
// 其他事件类型...
throw new IllegalArgumentException("Unknown event type");
}
}
应用服务
// 应用服务层
public class AccountApplicationService {
private final EventStore eventStore;
public AccountApplicationService(EventStore eventStore) {
this.eventStore = eventStore;
}
// 创建账户
public void createAccount(String accountId, String holderName) {
BankAccount account = BankAccount.createAccount(accountId, holderName);
saveEvents(account);
}
// 存款
public void deposit(String accountId, BigDecimal amount) {
BankAccount account = loadAccount(accountId);
account.deposit(amount);
saveEvents(account);
}
// 取款
public void withdraw(String accountId, BigDecimal amount) {
BankAccount account = loadAccount(accountId);
account.withdraw(amount);
saveEvents(account);
}
// 冻结
public void freeze(String accountId, String reason) {
BankAccount account = loadAccount(accountId);
account.freeze(reason);
saveEvents(account);
}
// 查询账户快照
public BankAccountSnapshot getAccountSnapshot(String accountId) {
List<DomainEvent> events = eventStore.load(accountId);
BankAccount account = BankAccount.rebuild(
events.stream()
.filter(e -> e instanceof AccountEvent)
.map(e -> (AccountEvent) e)
.collect(Collectors.toList())
);
return new BankAccountSnapshot(
account.getAccountId(),
account.getAccountHolderName(),
account.getBalance(),
account.isFrozen(),
account.getVersion()
);
}
private BankAccount loadAccount(String accountId) {
List<DomainEvent> events = eventStore.load(accountId);
if (events.isEmpty()) {
throw new AccountNotFoundException("Account not found: " + accountId);
}
return BankAccount.rebuild(
events.stream()
.filter(e -> e instanceof AccountEvent)
.map(e -> (AccountEvent) e)
.collect(Collectors.toList())
);
}
private void saveEvents(BankAccount account) {
List<DomainEvent> pendingEvents = account.getPendingEvents();
eventStore.append(pendingEvents);
account.clearPendingEvents();
}
}
// 账户快照类
public class BankAccountSnapshot {
private final String accountId;
private final String accountHolderName;
private final BigDecimal balance;
private final boolean frozen;
private final int version;
public BankAccountSnapshot(String accountId, String accountHolderName,
BigDecimal balance, boolean frozen, int version) {
this.accountId = accountId;
this.accountHolderName = accountHolderName;
this.balance = balance;
this.frozen = frozen;
this.version = version;
}
// Getters...
}
测试示例
public class EventSourcingExample {
public static void main(String[] args) {
// 初始化基础设施
EventStore eventStore = new InMemoryEventStore();
AccountApplicationService service = new AccountApplicationService(eventStore);
try {
// 1. 创建账户
String accountId = "ACC-001";
service.createAccount(accountId, "张三");
System.out.println("账户创建成功!");
// 2. 存款操作
service.deposit(accountId, new BigDecimal("1000.00"));
System.out.println("存款 1000 元成功!");
// 3. 取款操作
service.withdraw(accountId, new BigDecimal("300.50"));
System.out.println("取款 300.50 元成功!");
// 4. 查询账户状态
BankAccountSnapshot snapshot = service.getAccountSnapshot(accountId);
System.out.println("账户余额: " + snapshot.getBalance());
System.out.println("账户版本: " + snapshot.getVersion());
System.out.println("账户持有人: " + snapshot.getAccountHolderName());
// 5. 冻结账户
service.freeze(accountId, "怀疑账户被盗");
System.out.println("账户已冻结");
// 6. 尝试在冻结状态下存款(应失败)
try {
service.deposit(accountId, new BigDecimal("100.00"));
} catch (IllegalStateException e) {
System.out.println("冻结状态下无法操作: " + e.getMessage());
}
// 7. 展示所有事件
System.out.println("\n========== 事件历史 ==========");
List<DomainEvent> events = eventStore.load(accountId);
for (DomainEvent event : events) {
System.out.println("事件版本: " + event.getVersion() + ", 时间: " + event.getOccurredAt());
if (event instanceof AccountCreatedEvent) {
AccountCreatedEvent e = (AccountCreatedEvent) event;
System.out.println(" - 创建账户: " + e.getAccountHolderName());
} else if (event instanceof MoneyDepositedEvent) {
MoneyDepositedEvent e = (MoneyDepositedEvent) event;
System.out.println(" - 存款: " + e.getAmount() + " 元");
} else if (event instanceof MoneyWithdrawnEvent) {
MoneyWithdrawnEvent e = (MoneyWithdrawnEvent) event;
System.out.println(" - 取款: " + e.getAmount() + " 元");
} else if (event instanceof AccountFrozenEvent) {
AccountFrozenEvent e = (AccountFrozenEvent) event;
System.out.println(" - 冻结原因: " + e.getReason());
}
}
} catch (Exception e) {
e.printStackTrace();
}
}
}
输出示例
账户创建成功!
存款 1000 元成功!
取款 300.50 元成功!
账户余额: 699.50
账户版本: 3
账户持有人: 张三
账户已冻结
冻结状态下无法操作: Account is frozen
========== 事件历史 ==========
事件版本: 1, 时间: 2024-01-15T10:30:00
- 创建账户: 张三
事件版本: 2, 时间: 2024-01-15T10:30:01
- 存款: 1000.00 元
事件版本: 3, 时间: 2024-01-15T10:30:02
- 取款: 300.50 元
事件版本: 4, 时间: 2024-01-15T10:30:03
- 冻结原因: 怀疑账户被盗
- 事件驱动:所有状态变更都通过事件表达
- 不可变事件:已发生的事件不可修改
- 聚合重建:通过事件流重建聚合状态
- 版本控制:每个事件都有版本号,支持乐观锁
- 完整审计:所有操作都有完整记录
- 时点查询:可以从任何时间点重建状态
这个案例展示了事件溯源的核心概念和实现方式,适合用于理解事件溯源的基本原理。