Java事件溯源案例

wen java案例 3

本文目录导读:

Java事件溯源案例

  1. 事件定义
  2. 聚合根
  3. 事件存储
  4. 应用服务
  5. 测试示例
  6. 输出示例

我将为您提供一个完整的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
  - 冻结原因: 怀疑账户被盗
  1. 事件驱动:所有状态变更都通过事件表达
  2. 不可变事件:已发生的事件不可修改
  3. 聚合重建:通过事件流重建聚合状态
  4. 版本控制:每个事件都有版本号,支持乐观锁
  5. 完整审计:所有操作都有完整记录
  6. 时点查询:可以从任何时间点重建状态

这个案例展示了事件溯源的核心概念和实现方式,适合用于理解事件溯源的基本原理。

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