事件溯源:这次怎么落地的

去年接手一个金融对账系统重构时,团队在数据库设计上争论了很久。

问题背景

这个系统要处理多平台的交易对账,核心需求有几个:

  • 必须能够追溯任何一笔交易在任意时间点的状态
  • 监管要求保留所有状态变更记录,不能直接 UPDATE
  • 需要支持历史数据回放和报表重算
  • 系统重构时不能丢数据

传统的 CRUD 模式做不到这些——你只看到当前状态,看不到怎么走到这一步的。事件溯源把每次状态变更都记录成不可变的事件,就像银行的交易流水。

核心概念:事件溯源到底在做什么

简单说,事件溯源就是把"当前状态"换成"事件流"。

传统方式:

  • 订单表存当前状态:status=PAID
  • 每次状态变更 UPDATE 一行

事件溯源方式:

  • 订单表只存初始状态
  • 每次状态变更追加一条事件:OrderPaid, OrderShipped, OrderCompleted
  • 当前状态通过回放所有事件计算出来

存储方式变了,读写路径、一致性假设和团队分工也得跟着改。

事件定义:从实际问题出发

事件定义有几个实际要注意的点。我们一开始把事件设计得太细,结果后续维护变得很复杂。

这个是我们的最终事件结构:

{
  "eventId": "evt_8a3d9f4e2b1c",
  "aggregateId": "order_123456",
  "aggregateType": "Order",
  "eventType": "OrderPaid",
  "eventData": {
    "amount": 299.00,
    "paymentMethod": "wechat",
    "transactionId": "wxpay_20250101223344"
  },
  "eventTypeVersion": 1,
  "timestamp": "2025-01-01T22:33:44Z",
  "correlationId": "req_abc123",
  "causationId": "req_def456",
  "userId": "user_789",
  "metadata": {
    "source": "api",
    "ip": "192.168.1.100"
  }
}

这里有几个坑需要说清楚:

事件粒度:我们一开始把 OrderPaymentStartedOrderPaymentProcessingOrderPaymentCompleted 分成了三个事件,结果回放时一个付款操作要触发三次状态转换,性能很差。后来合并成一个 OrderPaid 事件,把中间状态放在 metadata 里。

事件数据: eventData 里只放业务相关的数据,不要把技术细节塞进去。比如我们最初把 retryCountprocessingTime 也放进 eventData,后来发现这些对业务没有价值,应该在 metadata 里。

版本控制:事件结构会变化,必须加版本号。我们改过一次 OrderCreated 事件,加了 customerId 字段,但因为没改版本号,导致旧事件回放失败。

事件存储:实际落地的选择

事件存储的选择直接关系到系统可靠性。我们评估了三种方案:

1. 传统关系型数据库

最简单的是用一张大表存所有事件:

CREATE TABLE events (
  event_id VARCHAR(64) PRIMARY KEY,
  aggregate_id VARCHAR(64) NOT NULL,
  aggregate_type VARCHAR(32) NOT NULL,
  event_type VARCHAR(64) NOT NULL,
  event_data JSON NOT NULL,
  event_type_version INT NOT NULL DEFAULT 1,
  timestamp TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
  correlation_id VARCHAR(64),
  causation_id VARCHAR(64),
  user_id VARCHAR(64),
  metadata JSON,
  INDEX idx_aggregate_id (aggregate_id, timestamp),
  INDEX idx_event_type (event_type, timestamp),
  INDEX idx_timestamp (timestamp)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

这种方式的好处是简单、好理解,坏处是写入量大时性能瓶颈明显。我们的系统高峰期每秒要写入 2000+ 事件,单表扛不住了。

2. 专用事件存储

我们试过 Axon Server,部署简单,开箱即用,但有几个实际问题:

  • 运维团队不熟悉,出问题要等供应商
  • 和现有基础设施集成不够平滑
  • 监控和告警体系要重新建立

最后放弃了——功能够,但运维团队不熟,和现有基础设施也接不顺,未知风险太多。

3. 消息队列 + 持久化层

这是最终的方案,用 Kafka 做事件总线,MySQL 做持久化层:

public class EventStore {
    private final KafkaTemplate<String, String> kafkaTemplate;
    private final EventRepository eventRepository;

    @Transactional
    public void save(DomainEvent event) {
        // 先写入数据库
        EventEntity entity = new EventEntity(event);
        eventRepository.save(entity);

        // 再发送到 Kafka
        kafkaTemplate.send(
            event.getAggregateType() + ".events",
            event.getAggregateId(),
            JsonUtils.toJson(event)
        );
    }
}

这种方式的好处是:

  • Kafka 提供了天然的持久化和重放能力
  • 写入性能高,能抗住高峰流量
  • 可以多消费者并行处理不同类型的事件

但要注意顺序性问题。同一个 aggregate 的事件必须顺序写入,我们通过 partition key 来保证:

kafkaTemplate.send(
    topic,
    event.getAggregateId(), // 用 aggregateId 作为 partition key
    JsonUtils.toJson(event)
);

事件回放:性能优化实战

事件回放是事件溯源的核心能力,也是最容易出性能问题的地方。

基础回放实现

最简单的回放逻辑是查询一个 aggregate 的所有事件,然后顺序重放:

public Order replayOrderEvents(String orderId) {
    List<DomainEvent> events = eventRepository
        .findByAggregateIdOrderByTimestampAsc(orderId);

    Order order = new Order();
    for (DomainEvent event : events) {
        order.apply(event);
    }
    return order;
}

这在事件数量少的时候没问题,但一个订单如果有几百个事件,性能就很差了。

快照策略

我们在事件数量超过 100 个时生成快照:

public Order replayOrderWithSnapshot(String orderId) {
    // 先找最新的快照
    Snapshot snapshot = snapshotRepository
        .findLatestByAggregateId(orderId);

    Order order = snapshot != null
        ? snapshot.restoreOrder()
        : new Order();

    // 只回放快照之后的事件
    long lastEventVersion = snapshot != null
        ? snapshot.getVersion()
        : 0;

    List<DomainEvent> events = eventRepository
        .findByAggregateIdAndVersionAfter(
            orderId,
            lastEventVersion
        );

    for (DomainEvent event : events) {
        order.apply(event);
    }

    return order;
}

快照生成策略要权衡:

  • 事件数量:我们选 100 个事件生成一次
  • 时间间隔:每天凌晨生成一次
  • 手动触发:重要操作后立即生成

快照结构:

{
  "aggregateId": "order_123456",
  "aggregateType": "Order",
  "version": 150,
  "snapshotData": {
    "orderId": "order_123456",
    "status": "COMPLETED",
    "amount": 299.00,
    "customerId": "customer_789",
    "shippingAddress": "上海市浦东新区..."
  },
  "timestamp": "2025-01-15T10:00:00Z"
}

并行回放优化

批量回放时可以用多线程:

public List<Order> replayOrdersBatch(List<String> orderIds) {
    return orderIds.parallelStream()
        .map(this::replayOrderWithSnapshot)
        .collect(Collectors.toList());
}

但要注意资源占用,我们控制并发数为 CPU 核心数 * 2:

System.setProperty(
    "java.util.concurrent.ForkJoinPool.common.parallelism",
    String.valueOf(Runtime.getRuntime().availableProcessors() * 2)
);

CQRS:读写分离实践

事件溯源天然适合和 CQRS 结合,因为写操作关注事件,读操作关注聚合状态。

读模型设计

读模型可以直接用传统的 SQL 查询,因为不需要保证一致性:

-- 订单查询表
CREATE TABLE order_read_model (
    order_id VARCHAR(64) PRIMARY KEY,
    status VARCHAR(32) NOT NULL,
    amount DECIMAL(10, 2) NOT NULL,
    customer_id VARCHAR(64),
    customer_name VARCHAR(128),
    created_at TIMESTAMP NOT NULL,
    updated_at TIMESTAMP NOT NULL,
    INDEX idx_status (status),
    INDEX idx_customer (customer_id),
    INDEX idx_created_at (created_at)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

事件投影

通过事件处理器更新读模型:

@Component
public class OrderProjector {
    private final OrderReadModelRepository readModelRepository;

    @EventHandler
    public void handle(OrderCreatedEvent event) {
        OrderReadModel model = new OrderReadModel();
        model.setOrderId(event.getOrderId());
        model.setStatus("CREATED");
        model.setAmount(event.getAmount());
        model.setCustomerId(event.getCustomerId());
        model.setCreatedAt(event.getTimestamp());
        model.setUpdatedAt(event.getTimestamp());
        readModelRepository.save(model);
    }

    @EventHandler
    public void handle(OrderPaidEvent event) {
        OrderReadModel model = readModelRepository
            .findById(event.getOrderId())
            .orElseThrow();

        model.setStatus("PAID");
        model.setUpdatedAt(event.getTimestamp());
        readModelRepository.save(model);
    }
}

最终一致性

读模型是最终一致的,这需要和业务方沟通好。我们的对账系统可以接受最多 5 秒的数据延迟,所以用了异步事件处理:

@Async("eventHandlerExecutor")
public void handleEvent(DomainEvent event) {
    // 异步处理事件,不阻塞主流程
    orderProjector.handle(event);
}

实际踩过的坑

1. 事件丢失

事件从数据库到 Kafka 的过程中有可能会丢,我们用了两步保障:

第一重保障:数据库事务

@Transactional
public void saveEvent(DomainEvent event) {
    // 事件和业务操作在同一个事务里
    eventRepository.save(new EventEntity(event));
    // 业务操作...
}

第二重保障:消息确认

@KafkaListener(topics = "order.events")
public void processEvent(ConsumerRecord<String, String> record) {
    try {
        DomainEvent event = JsonUtils.fromJson(
            record.value(),
            DomainEvent.class
        );
        orderProjector.handle(event);
    } catch (Exception e) {
        // 处理失败,不确认消息,触发重试
        throw e;
    }
}

2. 事件顺序混乱

同一个 aggregate 的事件必须顺序处理,但 Kafka 不能绝对保证顺序。我们的解决方案是:

public class EventProcessor {
    private final Map<String, BlockingQueue<DomainEvent>> pendingEvents
        = new ConcurrentHashMap<>();

    public void processEvent(DomainEvent event) {
        String aggregateId = event.getAggregateId();

        pendingEvents
            .computeIfAbsent(aggregateId, k -> new LinkedBlockingQueue<>())
            .add(event);

        processPendingEvents(aggregateId);
    }

    private void processPendingEvents(String aggregateId) {
        BlockingQueue<DomainEvent> queue = pendingEvents.get(aggregateId);
        while (!queue.isEmpty()) {
            DomainEvent event = queue.peek();
            if (canProcessEvent(event)) {
                queue.poll();
                doProcessEvent(event);
            } else {
                break; // 不能处理,等待下次
            }
        }
    }

    private boolean canProcessEvent(DomainEvent event) {
        // 检查是否是期望的下一个事件
        // 通过版本号或时间戳判断
    }
}

3. 快照不一致

快照和事件不一致会导致数据错乱,我们的做法是:

public void createSnapshot(Order order, long version) {
    // 先验证快照版本和事件版本是否一致
    long currentEventVersion = eventRepository
        .findMaxVersionByAggregateId(order.getOrderId());

    if (currentEventVersion != version) {
        throw new SnapshotVersionMismatchException(
            "Snapshot version " + version +
            " does not match current event version " +
            currentEventVersion
        );
    }

    snapshotRepository.save(new Snapshot(order, version));
}

4. 事件 Schema 演进

事件结构变化是必然的,我们用了适配器模式:

public interface EventAdapter<T extends DomainEvent> {
    T adapt(JsonNode eventData, int version);
}

public class OrderPaidEventAdapter
    implements EventAdapter<OrderPaidEvent> {

    @Override
    public OrderPaidEvent adapt(JsonNode eventData, int version) {
        if (version == 1) {
            return adaptFromV1(eventData);
        } else if (version == 2) {
            return adaptFromV2(eventData);
        }
        throw new UnsupportedEventVersionException(version);
    }

    private OrderPaidEvent adaptFromV1(JsonNode eventData) {
        // 处理 V1 版本的事件
        // 转换成当前版本的结构
    }

    private OrderPaidEvent adaptFromV2(JsonNode eventData) {
        // 处理 V2 版本的事件
    }
}

适用场景判断

不是所有系统都适合用事件溯源,我们总结了几个判断条件:

适合:

  • 需要完整审计轨迹的系统(金融、医疗、政务)
  • 复杂业务规则,需要追溯决策过程
  • 需要历史数据回放和重算
  • 业务变化频繁,需要灵活性

不适合:

  • 简单 CRUD 应用,用传统 CRUD 就够了
  • 对强一致性要求极高的系统
  • 团队不熟悉事件驱动模式,学习成本高
  • 事件量小,用事件溯源反而增加复杂度

收个尾

对账系统切到事件溯源后,审计轨迹、历史回放、报表重算这几块省了不少额外开发;代价是开发复杂度上去、团队要重新熟悉运维监控。

监管要留痕、状态变更不能原地 UPDATE、还要能回放重算——这几个约束叠在一起,事件溯源算是对症。简单 CRUD、强一致实时读、团队没事件驱动经验,我会慎重。

想试的话,先挑一个小模块跑通,别一上来动核心链路。技术选型对了,团队接不住,照样翻车。

版权声明: 本文首发于 指尖魔法屋-事件溯源:这次怎么落地的https://blog.thinkmoon.cn/post/127-event-sourcing-guide/) 转载或引用必须申明原指尖魔法屋来源及源地址!