事件溯源:这次怎么落地的
去年接手一个金融对账系统重构时,团队在数据库设计上争论了很久。
问题背景
这个系统要处理多平台的交易对账,核心需求有几个:
- 必须能够追溯任何一笔交易在任意时间点的状态
- 监管要求保留所有状态变更记录,不能直接 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"
}
}
这里有几个坑需要说清楚:
事件粒度:我们一开始把 OrderPaymentStarted、OrderPaymentProcessing、OrderPaymentCompleted 分成了三个事件,结果回放时一个付款操作要触发三次状态转换,性能很差。后来合并成一个 OrderPaid 事件,把中间状态放在 metadata 里。
事件数据: eventData 里只放业务相关的数据,不要把技术细节塞进去。比如我们最初把 retryCount、processingTime 也放进 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/) 转载或引用必须申明原指尖魔法屋来源及源地址!
评论
使用 GitHub 账号登录后即可留言,支持 Markdown。