事件驱动架构:用异步解耦微服务

·... 次阅读
分享:

为什么需要事件驱动

同步调用(HTTP/gRPC)的问题是调用链耦合

用户下单 → 订单服务 → 库存服务 → 支付服务 → 物流服务

任何一环超时或失败,整个链路中断。事件驱动改变了这种模式:

用户下单 → 订单服务 → 发布"订单已创建"事件

      ┌───────┼───────┐
      ↓       ↓       ↓
   库存服务  支付服务  通知服务

订单服务只管发布事件,谁关心这个事件谁就订阅。生产者不关心消费者,消费者不依赖生产者。

核心概念

事件 vs 命令

  • 事件(Event):已经发生的事实,过去式。如 OrderCreatedPaymentCompleted
  • 命令(Command):要求做某事的指令。如 CreateOrderProcessPayment

事件驱动架构中,服务通过事件通信,而非直接发命令。

消息队列 vs 事件总线

消息队列(Queue) 事件总线(Topic/Exchange)
消费模式 点对点(一条消息一个消费者) 发布-订阅(一条消息多个消费者)
典型产品 RabbitMQ(Queue 模式) Kafka、RabbitMQ(Topic 模式)
适用场景 任务分发、削峰填谷 事件广播、数据同步

实践方案

方案一:Kafka 事件流

Kafka 的高吞吐、持久化和消息重放特性,非常适合事件驱动:

订单服务 → Kafka Topic: order-events

    ┌───────────┼───────────┐
    ↓           ↓           ↓
 库存消费者   支付消费者   通知消费者

关键设计

  • Topic 按业务领域划分(order-events、payment-events、user-events)
  • 消息包含 eventIdeventTypetimestampaggregateIdpayload
  • 消费者通过 aggregateId 做幂等处理

消息结构示例

{
  "eventId": "evt_abc123",
  "eventType": "OrderCreated",
  "timestamp": "2026-07-12T10:30:00Z",
  "aggregateId": "order_456",
  "payload": {
    "userId": "user_789",
    "items": [
      { "productId": "prod_1", "quantity": 2 }
    ],
    "totalAmount": 199.00
  }
}

方案二:RabbitMQ + 发布订阅

RabbitMQ 的 Exchange + Queue 模型更灵活:

# 发布事件
channel.basic_publish(
    exchange='microservice_events',
    routing_key='order.created',
    body=json.dumps(event_data)
)

# 消费事件
channel.queue_bind(
    queue='inventory_service_queue',
    exchange='microservice_events',
    routing_key='order.created'
)

事件驱动的常见模式

1. Saga 模式(编排)

每个服务处理完本地事务后,发布事件触发下一个服务:

订单服务(创建订单) → [OrderCreated] → 库存服务(扣库存) → [InventoryReserved] → 支付服务(扣款)

如果某步失败,通过补偿事件回滚已执行的操作:

支付失败 → [PaymentFailed] → 库存服务(恢复库存) → [InventoryReleased] → 订单服务(取消订单)

2. CQRS + 事件溯源

将读写分离,通过事件重建状态:

写操作 → 发布事件 → 事件存储(Event Store) → 更新读模型(物化视图)

好处:

  • 读写可独立扩展
  • 完整的事件历史,可追溯、可审计
  • 可以随时重新构建读模型

3. 发件箱模式(Outbox Pattern)

避免“数据库写入成功但消息发送失败”的双写问题:

数据库事务 {
  INSERT INTO orders (...);
  INSERT INTO outbox (event_type, payload);
}
→ 定时任务读取 outbox → 发送到消息队列 → 标记为已发送

注意事项

  1. 幂等性:消费者可能重复消费,必须保证幂等(通过 eventId 去重)
  2. 消息顺序:同一聚合根的事件需要有序处理(Kafka 用 partition key)
  3. 最终一致性:接受短暂的数据不一致,设计补偿机制
  4. 监控与死信:处理失败的消息进入死信队列,需要告警和人工处理
  5. Schema 演化:事件格式向后兼容,推荐用 Protobuf 或 Avro 管理 Schema

同步 vs 异步:什么时候用事件驱动

场景 推荐方式
查询数据 同步(HTTP/gRPC)
实时性要求高(如支付扣款) 同步
数据同步、状态变更通知 异步事件
需要多个服务同时响应 异步事件
长流程业务(如审批流) 异步事件 + Saga
削峰填谷(如秒杀) 异步消息队列

评论