消息驱动架构(EDA)

事件驱动架构的核心概念、消息队列选型、事件契约设计与幂等性保障策略。

消息驱动架构(EDA)

事件驱动架构(Event-Driven Architecture)通过事件的发布与订阅实现服务解耦,是微服务异步协作的核心模式。

1. 核心模式

事件通知

Service A ──(Event)──→ Event Bus ──(Fan-out)──→ Service B
                                           ──→ Service C
                                           ──→ Service D

A 不感知 B/C/D,仅声明"发生了什么"

事件驱动状态转移(Event-Carried State Transfer)

事件携带完整状态数据,消费者无需回查生产者。

{
  "event_type": "OrderCreated",
  "order_id": "ORD-123",
  "customer_id": "C-456",
  "items": [...],
  "total": 299.00,
  "timestamp": "2026-08-13T10:00:00Z"
}

2. 消息队列选型

特性KafkaRabbitMQRocketMQPulsar
模式发布订阅 / 队列队列 / 路由发布订阅 / 队列发布订阅
吞吐量极高(百万级/s)中等(万级/s)高(十万级/s)极高
延迟ms 级微秒级ms 级ms 级
消息顺序Partition 内有序Queue 内有序Queue 内有序Partition 有序
持久化磁盘PageCache磁盘/Memory磁盘分层存储
重放能力支持 Offset 回溯消费者确认即删除支持支持
多租户支持 VHost支持 Namespace
场景日志/流处理任务队列/路由金融/电商云原生

Kafka 核心设计

Topic: events
├── Partition 0 → [E1][E2][E3] → Consumer Group A (C1, C2)
├── Partition 1 → [E4][E5][E6] → Consumer Group B
└── Partition 2 → [E7][E8]

Key Hash → Partition → 相同 Key 保证顺序
Offset → 消费者自行维护消费位点

3. 事件契约设计

Schema 管理

{
  "type": "record",
  "name": "OrderCreated",
  "namespace": "com.example.events",
  "fields": [
    { "name": "order_id", "type": "string" },
    { "name": "customer_id", "type": "string" },
    { "name": "total_amount", "type": "decimal", "precision": 10, "scale": 2 }
  ]
}

Schema 演进规则

操作向后兼容向前兼容说明
新增可选字段最安全
删除字段避免
修改字段类型新增字段替代
字段改名使用 aliases

Kafka Schema Registry(Confluent)

# 注册 Schema
curl -X POST http://schema-registry:8081/subjects/events-value/versions \
  -H "Content-Type: application/vnd.schemaregistry.v1+json" \
  -d '{"schema": "{...}"}'

# 序列化时自动注册/校验 Schema
producer.send(new ProducerRecord<>("events", avroRecord));

4. 幂等性设计

消息消费可能重复,消费者必须保证幂等。

幂等策略

策略适用场景
唯一键约束数据库插入(如订单号)
状态机校验业务状态流转控制
去重表 / 布隆过滤器消息级去重
乐观锁(版本号)并发更新

示例

@Transactional
public void onPaymentCompleted(PaymentCompletedEvent event) {
    // 方法一:唯一键
    try {
        orderRepository.markPaid(event.getOrderId(), event.getTransactionId());
    } catch (DuplicateKeyException e) {
        log.info("Duplicate payment event, ignored");
        return;
    }
    
    // 方法二:状态机
    Order order = orderRepository.findById(event.getOrderId());
    if (order.getStatus() != OrderStatus.PENDING_PAYMENT) {
        return; // 已处理过
    }
    order.pay(event.getTransactionId());
}

5. 事件溯源结合

┌──────────┐     Event          ┌──────────────┐
│ Command  │ ─────────────────→ │ Event Store  │
│ Handler  │                    │ (Event Log)  │
└──────────┘                    └──────┬───────┘
                                       │
                    Projections        │
                    ┌──────────────┐   │
                    │ Read Model   │ ←─┘
                    │ (CQRS)       │
                    └──────────────┘

6. 最佳实践

  1. 事件总线独立:核心业务事件与审计日志分开 Topic
  2. 至少一次投递:消息队列保证,消费者做幂等
  3. 死信队列:消费失败移至 DLQ,避免阻塞主队列
  4. 监控告警:堆积量、消费延迟、重试次数
  5. 事件版本v1.order.createdv2.order.created

总结

层面要点
选型按吞吐量、延迟、顺序需求选 Kafka/RabbitMQ/RocketMQ
契约Avro/Protobuf + Schema Registry
消费至少一次 + 幂等设计
运维监控堆积、死信、消费者滞后

EDA 是微服务解耦的关键手段,但引入了最终一致性复杂度,需权衡使用。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「架构」更多文章

  1. 架构评审与技术债管理
  2. SLA/SLO/SLI 与容量规划
  3. 云原生架构模式