事件驱动架构:从理念到落地
2026/6/11大约 8 分钟
事件驱动架构:从理念到落地
事件驱动架构:从理念到落地是系统设计的核心,它决定了系统的可扩展性、可靠性和可维护性。
本文介绍了事件驱动架构:从理念到落地的设计原则和实践经验,帮助你提升架构设计能力。
spring:
kafka:
producer:
bootstrap-servers: kafka-broker:9092
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
# 可靠性配置
acks: all # 所有 ISR 确认
retries: 3
enable-idempotence: true # 幂等生产者
compression-type: snappy # 压缩
### 4.3 RocketMQ:金融场景的首选
```java
// RocketMQ 事务消息 - 实现分布式事务
@Service
public class RocketMQOrderService {
@Autowired private RocketMQTemplate rocketMQTemplate;
@Transactional
public void createOrder(CreateOrderCommand command) {
Order order = orderRepository.save(command.toOrder());
// 事务消息:先发半消息,本地事务提交后消息才投递
rocketMQTemplate.sendMessageInTransaction(
"order-topic",
MessageBuilder.withPayload(OrderCreatedEvent.from(order)).build(),
order.getId() // 传递给事务检查器的参数
);
}
}
// 事务检查器:处理半消息的状态回查
@RocketMQTransactionListener
public class OrderTransactionListener
implements RocketMQLocalTransactionListener {
@Override
public RocketMQLocalTransactionState executeLocalTransaction(
Message msg, Object arg) {
try {
Long orderId = (Long) arg;
// 检查订单是否已成功创建
Order order = orderRepository.findById(new OrderId(orderId))
.orElse(null);
if (order != null) {
return RocketMQLocalTransactionState.COMMIT;
}
return RocketMQLocalTransactionState.ROLLBACK;
} catch (Exception e) {
return RocketMQLocalTransactionState.UNKNOWN;
}
}
@Override
public RocketMQLocalTransactionState checkLocalTransaction(
Message msg) {
// 回查逻辑:消息服务端定期回调此方法确认事务状态
Long orderId = extractOrderId(msg);
return orderRepository.existsById(new OrderId(orderId))
? RocketMQLocalTransactionState.COMMIT
: RocketMQLocalTransactionState.ROLLBACK;
}
}
// 消费者
@Service
@RocketMQMessageListener(
topic = "order-topic",
consumerGroup = "inventory-group",
consumeMode = ConsumeMode.ORDERLY // 顺序消费
)
public class InventoryConsumer
implements RocketMQListener<OrderCreatedEvent> {
@Override
public void onMessage(OrderCreatedEvent event) {
// RocketMQ 保证重试,需要实现幂等
if (inventoryDeductLogRepository.exists(event.getOrderId())) {
return; // 幂等:已处理过
}
inventoryService.deduct(event.getItems());
// 记录处理日志(用于幂等判断)
inventoryDeductLogRepository.save(
new DeductLog(event.getOrderId()));
}
}第五部分:实战——订单系统的事件驱动改造
5.1 改造前(同步耦合)
@Service
public class OrderService {
@Autowired private InventoryService inventoryService;
@Autowired private PaymentService paymentService;
@Autowired private NotificationService notificationService;
@Autowired private CouponService couponService;
@Autowired private LogService logService;
@Transactional
public Order createOrder(OrderRequest request) {
// 同步调用所有服务 → 耦合严重
inventoryService.deduct(request.getItems()); // 2s
couponService.use(request.getCouponId()); // 0.5s
Order order = saveOrder(request); // 0.1s
paymentService.createPayment(order); // 1s
notificationService.sendSms(order); // 0.5s
logService.recordAuditLog(order); // 0.2s
return order;
// 总耗时: 4.3 秒!用户体验极差
}
}5.2 改造后(事件驱动)
@Service
public class OrderService {
@Autowired private OrderRepository orderRepository;
@Autowired private DomainEventPublisher eventPublisher;
@Transactional
public Order createOrder(CreateOrderCommand command) {
// 1. 创建订单(核心操作)
Order order = Order.create(command);
orderRepository.save(order);
// 2. 发布领域事件
eventPublisher.publish(OrderCreatedEvent.from(order));
// 3. 立即返回
return order;
// 总耗时: ~0.1 秒
}
}
// === 各消费者独立处理 ===
// 库存服务 - 扣减库存
@Component
public class OrderInventoryHandler {
@EventListener
@Transactional
public void handleOrderCreated(OrderCreatedEvent event) {
for (OrderItem item : event.getItems()) {
inventoryRepository.deduct(
item.getProductId(), item.getQuantity());
}
}
}
// 支付服务 - 创建支付单
@Component
public class OrderPaymentHandler {
@EventListener
public void handleOrderCreated(OrderCreatedEvent event) {
paymentService.createPayment(
event.getOrderId(), event.getTotalAmount());
}
}
// 通知服务 - 发送短信
@Component
public class OrderNotificationHandler {
@EventListener
@Async // 异步处理
public void handleOrderCreated(OrderCreatedEvent event) {
notificationService.sendOrderConfirmation(
event.getCustomerId(), event.getOrderId());
}
}
// 审计日志服务
@Component
public class OrderAuditLogHandler {
@EventListener
@Async
public void handleOrderCreated(OrderCreatedEvent event) {
auditLogRepository.save(new AuditLog(
"ORDER_CREATED",
event.getOrderId(),
event.getCustomerId(),
Instant.now()
));
}
}5.3 事件版本管理
事件结构会随着业务发展而变化,需要做好版本管理:
/**
* 带版本号的订单创建事件
* 向后兼容:新字段增加默认值,旧字段不删除
*/
public class OrderCreatedEvent implements DomainEvent {
// 事件元数据
private String eventId;
private String eventType = "OrderCreated";
private int version; // 事件版本号
private Instant occurredAt;
// V1 字段
private OrderId orderId;
private CustomerId customerId;
private Money totalAmount;
// V2 新增字段(2025-03)
private String channel; // 下单渠道 (APP/WEB/MINI_PROGRAM)
private String promotionCode; // 推广码
// V3 新增字段(2025-06)
private boolean isPresale; // 是否预售
private Instant expectedShipDate; // 预计发货日期
// 构造方法
private OrderCreatedEvent() {} // 序列化用
public static OrderCreatedEvent from(Order order) {
OrderCreatedEvent event = new OrderCreatedEvent();
event.eventId = UUID.randomUUID().toString();
event.version = 3; // 当前版本
event.occurredAt = Instant.now();
event.orderId = order.getId();
event.customerId = order.getCustomerId();
event.totalAmount = order.getTotalAmount();
event.channel = order.getChannel();
event.promotionCode = order.getPromotionCode();
event.isPresale = order.isPresale();
event.expectedShipDate = order.getExpectedShipDate();
return event;
}
}
// 消费者做版本兼容处理
@Component
public class InventoryHandler {
@EventListener
public void handleOrderCreated(OrderCreatedEvent event) {
// 按版本处理不同逻辑
switch (event.getVersion()) {
case 1:
handleV1(event);
break;
case 2:
handleV2(event);
break;
case 3:
handleV3(event);
break;
default:
log.warn("未知事件版本: {}", event.getVersion());
// 尝试降级处理
handleV1(event);
}
}
private void handleV1(OrderCreatedEvent event) {
// 处理 V1 逻辑(没有渠道、预售等信息)
}
}事件版本管理最佳实践:
- 只增不减:新字段只添加,旧字段标记
@Deprecated但不删除 - 消费者容错:未知字段忽略(不报错)
- 版本号递增:每次结构变更,版本号 +1
- 降级兼容:新消费者应能处理旧版本事件
第六部分:最终一致性处理
6.1 拥抱最终一致性
事件驱动架构的关键转变:从强一致性到最终一致性。
强一致性(同步):
订单创建 → 等待库存扣减完成 → 返回给用户
用户等待时间 = 所有步骤耗时之和
最终一致性(异步):
订单创建 → 立即返回给用户
库存扣减(异步)→ 可能晚几秒完成
用户等待时间 = 订单创建耗时
代价:用户可能看到"订单已创建,库存处理中..."的状态
收益:系统吞吐量提升 10 倍以上6.2 处理失败重试
@Component
public class RobustInventoryHandler {
private static final int MAX_RETRIES = 3;
@EventListener
public void handleOrderCreated(OrderCreatedEvent event) {
for (int attempt = 1; attempt <= MAX_RETRIES; attempt++) {
try {
inventoryService.deduct(event.getItems());
return; // 成功
} catch (InsufficientInventoryException e) {
// 库存不足 - 不重试,直接处理
handleInsufficientInventory(event);
return;
} catch (TemporaryException e) {
// 临时错误 - 重试
log.warn("库存扣减临时失败,第 {} 次重试", attempt);
sleep(1000L * attempt); // 指数退避
}
}
// 重试耗尽 - 发送到死信队列
deadLetterQueue.send(event);
alertService.alert("库存扣减失败", event);
}
}6.3 幂等性保障
事件可能被重复投递(at-least-once 语义),消费者必须实现幂等:
@Component
public class IdempotentInventoryHandler {
@Autowired private EventProcessLogRepository logRepository;
@EventListener
@Transactional
public void handleOrderCreated(OrderCreatedEvent event) {
// 幂等检查:用 eventId 作为唯一键
String eventId = event.getEventId();
if (logRepository.existsByEventId(eventId)) {
log.info("事件已处理,跳过: eventId={}", eventId);
return; // 幂等:已处理过
}
// 处理业务
inventoryService.deduct(event.getItems());
// 记录处理日志(插入成功 = 唯一索引保证不会重复处理)
logRepository.save(new EventProcessLog(eventId, "PROCESSED", Instant.now()));
}
}
// 数据库表设计
CREATE TABLE t_event_process_log (
event_id VARCHAR(64) PRIMARY KEY, -- 唯一约束 = 幂等保证
status VARCHAR(20),
processed_at TIMESTAMP,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);第七部分:EDA 的陷阱
7.1 事件风暴
症状:一个事件触发另一个事件,形成事件链,最终无法追踪数据流向:
OrderCreated → InventoryDeducted → StockLowAlert
→ ReplenishmentRequested
→ ReplenishmentApproved
→ PurchaseOrderCreated
→ ...解决方案:
- 事件溯源记录所有事件(虽然不能阻止风暴,但能追溯)
- 明确的事件所有权:谁能发布什么事件
- 事件流图:用文档/图表显式化事件依赖关系
- 断路保护:检测事件循环并熔断
7.2 调试困难
症状:一个请求分散在多个服务的多个事件中,像"大海捞针":
排查"为什么用户没收到库存扣减通知":
1. 查订单服务的 order-created 事件是否发布
2. 查库存服务是否消费了该事件
3. 查 stock-low-alert 事件是否发布
4. 查通知服务是否消费了 stock-low-alert 事件
5. 查短信网关是否成功发送解决方案:
- 全链路追踪:Jaeger/Zipkin 关联所有事件
- 事件 ID 传递:在事件中携带 traceId
- 事件追踪平台:可视化事件流
- 每个事件携带 correlationId:串联同一个业务流程的所有事件
public class OrderCreatedEvent {
private String eventId;
private String traceId; // 链路追踪 ID(从上游提取)
private String correlationId; // 业务流程 ID(关联同一个流程)
// ...
}
// 发布事件时传递 traceId
public void publish(Order order, String traceId) {
OrderCreatedEvent event = OrderCreatedEvent.from(order);
event.setTraceId(traceId);
event.setCorrelationId(order.getOrderNo().toString());
eventBus.publish(event);
}7.3 数据不一致的排查
事件驱动下的数据不一致比同步调用更难排查,因为问题可能在你看到的时候已经"过去了"。
最佳实践:
- 监控事件延迟:如果事件处理延迟超过阈值,告警
- 对账机制:定期对比源系统和目标系统的数据
- 补偿任务:定期扫描"未完成"的业务流程并补偿
/**
* 每天凌晨 2 点执行的对账任务
*/
@Component
public class DailyReconciliationJob {
@Scheduled(cron = "0 0 2 * * ?")
public void reconcile() {
// 查询"订单已创建但库存未扣减"的异常数据
List<Order> unreconciledOrders = orderRepository
.findOrdersWithoutInventoryDeduction(
LocalDate.now().minusDays(1));
for (Order order : unreconciledOrders) {
log.warn("发现未对账订单: {}", order.getId());
// 补偿:重新发布库存扣减事件
eventPublisher.publish(
new InventoryDeductCompensationEvent(order));
}
}
}第八部分:架构演进路线
从单体到事件驱动三步走
Phase 1: 内部事件
在单体应用内使用 Spring Events (@EventListener)
订单模块事件 → 通知模块、日志模块
目标是理解事件思维,不引入外部依赖
Phase 2: 外部事件总线
引入 Kafka/RocketMQ
核心跨模块通信使用消息队列
非核心仍保留同步调用
Phase 3: 完整 EDA
所有模块间通信通过事件
全链路追踪
事件版本管理
自动化对账和补偿对于大多数项目,Phase 2 就是最优状态。Phase 3 适合事件驱动的产品型公司(如电商平台、金融交易系统)。
总结
| 概念 | 要点 |
|---|---|
| 事件通知 | 轻量,只发 ID,消费者回调查询 |
| 携带状态转移 | 推荐:事件自包含,消费者无需回调 |
| 事件溯源 | 不存状态存事件历史,适合审计 |
| Command vs Event | Command=请求,Event=事实 |
| 消息选型 | 大数据→Kafka,金融→RocketMQ,路由灵活→RabbitMQ |
| 幂等 | 靠 eventId + 唯一索引保证 |
| 版本管理 | 只增不减,消费者容错 |
| 对账 | 定期检查数据一致性,补偿异常 |
最终一致性不是 bug,是 feature。 拥抱它,用对账和监控来管理它,而不是试图消灭它。
参考资料
- Martin Fowler, "What do you mean by 'Event-Driven'?"
- Greg Young, "CQRS and Event Sourcing"
- Confluent, "Event-Carried State Transfer"
- Chris Richardson, "Pattern: Event-driven architecture"
- RocketMQ 官方文档 - 事务消息