消息驱动架构(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. 消息队列选型
| 特性 | Kafka | RabbitMQ | RocketMQ | Pulsar |
|---|---|---|---|---|
| 模式 | 发布订阅 / 队列 | 队列 / 路由 | 发布订阅 / 队列 | 发布订阅 |
| 吞吐量 | 极高(百万级/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. 最佳实践
- 事件总线独立:核心业务事件与审计日志分开 Topic
- 至少一次投递:消息队列保证,消费者做幂等
- 死信队列:消费失败移至 DLQ,避免阻塞主队列
- 监控告警:堆积量、消费延迟、重试次数
- 事件版本:
v1.order.created→v2.order.created
总结
| 层面 | 要点 |
|---|---|
| 选型 | 按吞吐量、延迟、顺序需求选 Kafka/RabbitMQ/RocketMQ |
| 契约 | Avro/Protobuf + Schema Registry |
| 消费 | 至少一次 + 幂等设计 |
| 运维 | 监控堆积、死信、消费者滞后 |
EDA 是微服务解耦的关键手段,但引入了最终一致性复杂度,需权衡使用。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。