Skip to content
第 220 / 250 章架构⏱ 12 分钟阅读

第 220 章:事件驱动架构

学习目标

  • 理解事件驱动架构(EDA)
  • 掌握 Event Sourcing 模式
  • 学会 CQRS 模式
  • 了解 EventStorming 工作坊

一、事件驱动简介

事件驱动架构(Event-Driven Architecture)以事件为核心,服务间通过发布/订阅事件通信。

1.1 同步 vs 异步通信

维度同步(RPC)异步(事件)
调用方式请求-响应发布-订阅
一致性强一致最终一致
性能受最慢接口高吞吐
耦合知道对方不关心对方
失败处理即时失败重试 + DLQ
适用核心路径副作用、通知

1.2 事件驱动的好处

  1. 解耦:服务间不直接调用
  2. 可扩展:新订阅者无需改发布者
  3. 可观测:事件流可追溯
  4. 容错:重试、DLQ 保留
  5. 演进:加新功能只需加订阅者

二、事件基础

2.1 事件 vs 命令

维度事件(Event)命令(Command)
时态过去式(已发生)祈使句(请执行)
OrderPaidPayOrder
发起方聚合根用户 / 外部
处理多个订阅者一个处理器

2.2 事件结构

java
public interface DomainEvent {
    String eventId();         // 全局唯一
    Instant occurredOn();     // 发生时间
    String eventType();       // 类型,如 "order.paid"
    String aggregateId();     // 聚合根 ID
}

具体事件:

java
public record OrderPaidEvent(
    String eventId,
    Instant occurredOn,
    String orderId,
    Money amount,
    String customerId
) implements DomainEvent {
}

2.3 事件命名

过去式,业务含义:

✅ OrderCreated / OrderPaid / OrderShipped
❌ CreateOrder / PayOrder  ← 这是命令

三、事件发布

3.1 Spring ApplicationEvent

java
@Service
public class OrderService {

    @Autowired
    private ApplicationEventPublisher publisher;

    @Transactional
    public void pay(OrderId id) {
        Order order = orderRepo.findById(id);
        order.pay();
        orderRepo.save(order);

        // 同步发布(Spring 事务内)
        publisher.publishEvent(new OrderPaidEvent(id));
    }
}

@Component
public class OrderPaidListener {
    @EventListener
    public void handle(OrderPaidEvent e) {
        log.info("订单已支付: {}", e.orderId());
    }
}

3.2 通过消息队列

java
@Service
@RequiredArgsConstructor
public class OrderEventPublisher {
    private final KafkaTemplate<String, DomainEvent> kafka;

    public void publish(DomainEvent event) {
        // 事务性发布:绑定消息与本地事务(见 Outbox 模式)
        kafka.send("order-events", event.eventType(), event);
    }
}

3.3 Outbox 模式(发布保证)

问题:数据库 commit + MQ publish,可能数据 commit 而消息丢失。

sql
-- outbox 表
CREATE TABLE outbox (
    id BIGSERIAL PRIMARY KEY,
    aggregate_id TEXT,
    event_type TEXT,
    payload JSONB,
    created_at TIMESTAMP,
    sent_at TIMESTAMP NULL
);
java
@Transactional
public void createOrder() {
    Order order = ...;
    orderRepo.save(order);
    // 同一事务插入 outbox
    outboxRepo.save(new OutboxEntry(
        order.id(),
        "OrderCreated",
        objectMapper.writeValueAsString(event)
    ));
}

@Scheduled(fixedDelay = 1000)
public void pollOutbox() {
    outboxRepo.findUnsent().forEach(entry -> {
        kafka.send(entry.eventType(), entry.payload);
        entry.markSent();
    });
}

四、事件订阅

4.1 同步事件订阅

java
@Service
public class InventoryService {
    private final Cache<String, Boolean> stockCache;

    @EventListener
    public void on(OrderConfirmedEvent e) {
        e.items().forEach(item ->
            stockCache.put(item.skuId(), false)   // 已锁库存
        );
    }
}

4.2 异步消息订阅

java
@Component
public class InventoryEventSubscriber {

    @KafkaListener(topics = "order-events", groupId = "inventory")
    public void handle(OrderConfirmedEvent event) {
        event.items().forEach(item ->
            inventoryRepo.lockStock(item.skuId(), item.quantity())
        );
    }
}

4.3 幂等消费

At-Least-Once + 幂等:

java
@KafkaListener(topics = "order-events")
public void handle(OrderConfirmedEvent event) {
    // 1. 检查是否已处理
    if (processedRepo.exists(event.eventId())) {
        return;
    }

    // 2. 处理业务
    inventoryService.lockStock(event);

    // 3. 标记已处理
    processedRepo.save(event.eventId());
}

数据库唯一约束:

sql
CREATE UNIQUE INDEX uk_processed_event ON processed_events(event_id);

五、CQRS 模式

5.1 概念

Command Query Responsibility Segregation

  • 写模型:聚合根,完整业务规则
  • 读模型:为查询优化(扁平、冗余)
  • 通过事件同步

5.2 为什么需要 CQRS

场景CRUDCQRS
简单业务过度
读写比例 100:1,读慢
多查询模型
高一致性复杂
复杂业务(订单、库存)

5.3 示例:订单查询

写模型(Command 侧):

java
public class Order {
    private OrderId id;
    private CustomerId customerId;
    private List<OrderItem> items;
    private Money total;
    private OrderStatus status;
    // ... 业务行为
}

读模型(Query 侧):

java
@Data
public class OrderView {
    private String orderId;
    private String customerName;        // 冗余查询高效
    private String customerPhone;       // 冗余
    private List<String> skuNames;      // 冗余
    private BigDecimal total;
    private String statusName;
    private Instant createdAt;
}

事件处理器构建读模型:

java
@KafkaListener(topics = "order-events", groupId = "query-side")
public void project(DomainEvent event) {
    if (event instanceof OrderCreatedEvent e) {
        var view = new OrderView();
        view.setOrderId(e.orderId());
        view.setCustomerName(userService.getName(e.customerId()));
        view.setStatusName("待支付");
        orderViewRepo.save(view);
    }

    if (event instanceof OrderPaidEvent e) {
        orderViewRepo.updateStatus(e.orderId(), "已支付");
    }
}

查询:

java
@GetMapping("/orders/{id}/detail")
public OrderView detail(@PathVariable String id) {
    return orderViewRepo.findById(id);  // 单表,快
}

六、Event Sourcing(事件溯源)

6.1 概念

不存储对象的当前状态,只存储导致状态变化的事件流

6.2 实现

java
public class Order {
    private OrderId id;
    private List<OrderItem> items = new ArrayList<>();
    private Money total;
    private OrderStatus status;

    // 应用事件,改变状态
    public static Order rehydrate(List<DomainEvent> events) {
        Order order = new Order();
        events.forEach(order::apply);
        return order;
    }

    private void apply(DomainEvent event) {
        if (event instanceof OrderCreatedEvent e) {
            this.id = new OrderId(e.orderId());
            this.status = OrderStatus.PENDING;
        } else if (event instanceof ItemAddedEvent e) {
            this.items.add(new OrderItem(e.productId(), e.qty()));
        } else if (event instanceof OrderPaidEvent e) {
            this.status = OrderStatus.PAID;
        }
        // ...
    }

    // 命令:产生事件
    public void addItem(ProductId productId, int qty) {
        ItemAddedEvent event = new ItemAddedEvent(id.value(), productId.value(), qty);
        apply(event);
        registerEvent(event);
    }
}

6.3 存储

sql
-- 事件存储
CREATE TABLE event_store (
    id BIGSERIAL PRIMARY KEY,
    aggregate_id TEXT,
    aggregate_type TEXT,
    event_type TEXT,
    event_data JSONB,
    version INT,
    created_at TIMESTAMP
);

CREATE INDEX idx_aggregate ON event_store(aggregate_id, version);
java
public class EventStoreRepository {
    public List<DomainEvent> load(String aggregateId) {
        return jdbc.query(
            "SELECT * FROM event_store WHERE aggregate_id = ? ORDER BY version",
            eventRowMapper,
            aggregateId
        );
    }

    public void save(Order order) {
        var version = currentVersion(order.id());
        order.getEvents().forEach(event -> {
            jdbc.update(
                "INSERT INTO event_store(aggregate_id, event_type, event_data, version) VALUES (?, ?, ?, ?)",
                order.id().value(), event.getClass().getSimpleName(),
                toJson(event), ++version
            );
        });
    }
}

6.4 快照(Snapshot)

事件量大时,定时保存快照 + 仅回放快照后的事件。

java
public class OrderSnapshot {
    private String orderId;
    private List<OrderItem> items;
    private Money total;
    private OrderStatus status;
    private int version;
}
java
public Order load(String orderId) {
    var snapshot = snapshotRepo.findByOrder(orderId);
    var events = eventStore.loadFrom(orderId, snapshot.version());
    return Order.rehydrate(snapshot, events);
}

6.5 Event Sourcing 优劣

优势劣势
完整审计复杂,学习曲线
时序回放必须有快照机制
任意时点查询不适合简单 CRUD
天然支持 CQRS事件 schema 演进难

七、Saga 模式

跨服务事务拆为多个本地事务 + 补偿。

7.1 编排式 Saga

7.2 事件编排式(Choreography)

7.3 Seata Saga

阿里开源的分布式事务解决方案:

java
@GlobalTransactional(name = "createOrder", rollbackFor = Exception.class)
public void createOrder(OrderRequest req) {
    orderService.create(req);
    inventoryService.deduct(req);
    paymentService.pay(req);
}

八、事件溯源 + CQRS 实战

8.1 架构

8.2 命令处理

java
@RestController
@RequestMapping("/commands")
public class OrderCommandController {

    @Autowired
    private OrderCommandService service;

    @PostMapping("/orders")
    public ResponseEntity create(@RequestBody CreateOrderCmd cmd) {
        return ResponseEntity.ok(service.createOrder(cmd));
    }
}

@Service
public class OrderCommandService {

    @Autowired
    private EventStore eventStore;

    @Transactional
    public OrderId createOrder(CreateOrderCmd cmd) {
        Order order = new Order();
        order.create(cmd.customerId, cmd.items);

        eventStore.save(order);
        return order.id();
    }
}

8.3 投影处理

java
@KafkaListener(topics = "order-events", groupId = "read-model-projection")
public void project(DomainEvent event) {
    if (event instanceof OrderCreatedEvent e) {
        mongoTemplate.save(toDocument(e));
    } else if (event instanceof OrderItemAddedEvent e) {
        mongoTemplate.updateFirst(
            query(where("_id").is(e.orderId())),
            new Update().inc("itemCount", 1),
            "order_view"
        );
    }
}

8.4 查询

java
@QueryHandler
public OrderView handle(GetOrderQuery query) {
    // 直接读 MongoDB(查询专用)
    return mongoTemplate.findById(query.id(), OrderView.class);
}

九、事件演化

9.1 Schema 注册

每个事件有版本号,支持向上兼容。

json
{
  "eventType": "OrderPaid",
  "version": 1,
  "data": {
    "orderId": "123",
    "amount": "19.99",
    "paidAt": "2026-08-13T10:30:00Z"
  }
}

9.2 兼容策略

变更兼容
添加字段
添加事件类型
改字段语义❌ 发新版本
删除字段❌ 新事件发布

9.3 升级流程

  1. 发布 OrderPaid_v2(兼容 v1 数据)
  2. 老消费者读 v1,新消费者读 v2
  3. 全量迁移到 v2
  4. 弃用 v1

十、常见误区

误区教训
把事件当 RPC事件语义应"过去时",不可谓语
同步等待事件结果用 Saga 或反查
事件循环依赖A→B→A 会死锁
事件不幂等消费失败重试导致多次处理
Event 包含过多数据只放必要的状态变化
Event 缺少追溯信息至少含 eventId + occurredOn

十一、本章小结

模式适用
事件驱动服务解耦 + 异步通信
CQRS读写压力严重不均
Event Sourcing需完整审计 + 时序回放
Saga跨服务事务

动手练习

  1. 设计订单系统的事件流(Created / Paid / Shipped / Completed / Refunded)
  2. 用 Spring + Kafka 实现 Outbox 模式
  3. 实现一个简化的 CQRS:写模型 + 读模型 + 事件同步
  4. 设计一个 Saga:下单 → 库存锁定 → 支付 → 失败回滚

推荐阅读


下一章:第 221 章:Spring Cloud 入门

本站基于 VitePress 构建 · 由 Codebook 团队维护