上一篇【第89篇】实时数据同步平台的Kafka实战——MySQL CDC与Kafka的最佳组合
系列完结,感谢阅读


摘要

微服务架构走到深水区,你会发现最大的挑战不是"拆",而是拆完之后"怎么协作"。传统的同步RPC调用链让服务之间紧耦合——A调用B,B调用C,C调用D……任何一个环节慢,整个链路就慢。事件驱动架构用Kafka作为中枢神经系统,让服务之间通过"事件"而非"调用"来协作,彻底解耦。

本文是Kafka实战系列的收官之作,聚焦微服务架构中最核心的几个课题:事件驱动 vs 请求响应的根本差异、领域事件的正确设计姿势、Kafka作为事件总线的落地方案、分布式事务的Saga模式(Choreography编舞 vs Orchestration编排)选型指南,以及事件溯源Event Sourcing的实践心法。读完这篇,你对微服务架构的认知将从"能用"进阶到"用得好"。


一、事件驱动架构 vs 请求响应架构

1.1 两种架构的思维方式差异

【请求响应架构(RPC调用链)】

 ┌────────┐     ┌────────┐     ┌────────┐     ┌────────┐
 │ 订单服务 │────►│ 库存服务 │────►│ 支付服务 │────►│ 通知服务 │
 │ create  │     │ deduct  │     │ pay     │     │ notify  │
 └────────┘     └────────┘     └────────┘     └────────┘
     │               │               │               │
     同步等待        同步等待         同步等待          返回结果
     
  问题:
  • 调用链总超时 = 每个环节超时之和
  • 任何一环挂了,整个链路失败
  • 加一个新步骤?改订单服务代码


【事件驱动架构(Kafka事件总线)】

   ┌────────┐  发布事件  ┌──────────────────────────────┐
   │ 订单服务 │──────────►│      Kafka 事件总线           │
   │ create  │           │                               │
   └────────┘           │  Topic: order.events          │
                         │  ┌─────────────────────────┐ │
                         │  │ Event: order.created     │─┼──► 库存服务(订阅)
                         │  │ Event: order.paid        │─┼──► 支付服务(订阅)
                         │  │ Event: order.shipped     │─┼──► 物流服务(订阅)
                         │  │ Event: order.delivered   │─┼──► 通知服务(订阅)
                         │  └─────────────────────────┘ │
                         └──────────────────────────────┘

  优点:
  • 订单服务不关心谁消费——只管发事件
  • 新服务加入零侵入——订阅即可
  • 各服务独立处理,互不阻塞
  • 天然支持事件回溯和审计

1.2 全面对比

维度请求响应架构事件驱动架构
耦合度高(调用方知道所有下游)低(只关心事件,不关心消费者)
扩展性困难(改调用方代码)简单(新服务订阅即可)
可用性取决于最慢的一环各服务独立,相互隔离
延迟低(直接调用)稍高(经过Kafka中转)
一致性强一致性(同步)最终一致性(异步)
可观测性需要链路追踪事件流=天然的审计日志
复杂度简单直接需要考虑消息顺序、重复、幂等
适用场景简单CRUD、要求强一致复杂业务流程、多下游消费

什么时候选事件驱动? 当你的业务流程涉及3个以上下游服务,或者需要频繁新增下游时,事件驱动就是正确的选择。


二、领域事件设计原则

2.1 好的事件设计

事件不是"数据库变更通知",而是"业务事实的不可变记录"。

【不良设计 vs 良好设计】

不良设计(技术事件):
  ❌ "orders表的status字段从CREATED变成PAID"
  ❌ "user_point加100"  
  ❌ "第3行被更新了"

良好设计(领域事件):
  ✅ "订单已支付"  → order.paid
  ✅ "用户完成新手任务,获得100积分"  → task.completed
  ✅ "物流信息更新,订单已发货"  → order.shipped

2.2 事件设计原则

原则说明示例
不可变事件发布后不可修改用eventId+timestamp标识
过去时命名用过去时态order.created ✅ / order.create ❌
自描述事件包含处理所需的所有数据包含orderId、userId、amount等
瘦事件只包含标识符和关键数据orderId而不是整个Order对象
版本化事件格式要支持演进用version字段

2.3 事件Payload设计

// 好的事件设计:自描述 + 瘦事件
@Data
@Builder
public class OrderPaidEvent {
    // ===== 事件元数据 =====
    private String eventId;           // 事件唯一ID
    private String eventType;         // "order.paid"
    private long timestamp;           // 发生时间
    private String version;           // 事件格式版本 "v1"
    private String correlationId;     // 关联ID(链路追踪)
    
    // ===== 业务关键数据 =====
    private String orderId;           // 订单ID
    private String userId;            // 用户ID
    private long paidAmount;          // 支付金额(分)
    private String paymentMethod;     // 支付方式
    
    // ===== 注意:不要放整个Order对象! =====
    // 消费者如果需要更多信息,自己去查数据库
}

// 不好的事件设计
@Data
class BadEvent {
    // ❌ 没有eventId → 无法去重
    // ❌ 命名为paid而是过去式 → 语义不清晰  
    String type;  // "paid"
    // ❌ 放了整个订单对象 → 臃肿 + 容易暴露敏感数据
    FullOrder order;
}

三、Kafka作为事件总线——落地方案

3.1 事件总线架构

【Kafka事件总线架构】

 ┌────────────────────────────────────────────────────────────┐
 │                   Kafka 事件总线                             │
 │                                                              │
 │  ┌──────────────────────────────────────────────────────┐  │
 │  │              Domain Events (领域事件)                 │  │
 │  │                                                       │  │
 │  │  order.events.v1      payment.events.v1              │  │
 │  │  user.events.v1       inventory.events.v1            │  │
 │  │  notification.events.v1                              │  │
 │  └──────────────────────────────────────────────────────┘  │
 │                                                              │
 │  ┌──────────────────────────────────────────────────────┐  │
 │  │         Integration Events (集成事件/跨边界)           │  │
 │  │                                                       │  │
 │  │  delivery.external.carrier.events                    │  │
 │  │  sms.external.vendor.events                          │  │
 │  └──────────────────────────────────────────────────────┘  │
 └────────────────────────────────────────────────────────────┘

区分领域事件和集成事件:
  • 领域事件: 内部服务间通信,格式可控
  • 集成事件: 与外部系统对接,需要适配转换

3.2 Spring Cloud Stream集成

// ===== 生产者端 =====
@Component
public class OrderEventPublisher {
    
    private final StreamBridge streamBridge;
    
    public void publishOrderCreated(Order order) {
        OrderCreatedEvent event = OrderCreatedEvent.builder()
            .eventId(UUID.randomUUID().toString())
            .eventType("order.created")
            .timestamp(System.currentTimeMillis())
            .version("v1")
            .orderId(order.getId())
            .userId(order.getUserId())
            .amount(order.getAmount())
            .build();
        
        // 发送到Kafka
        Message<OrderCreatedEvent> message = MessageBuilder
            .withPayload(event)
            .setHeader("eventType", "order.created")
            .setHeader("messageKey", order.getId())  // 分区Key
            .build();
        
        streamBridge.send("orderEvents-out-0", message);
    }
}

// ===== 消费者端 =====
@Component
public class InventoryEventHandler {
    
    private final InventoryService inventoryService;
    
    @Bean
    public Consumer<Message<OrderCreatedEvent>> processOrderCreated() {
        return message -> {
            OrderCreatedEvent event = message.getPayload();
            
            // 幂等检查(根据eventId去重)
            if (eventDedupService.isDuplicate(event.getEventId())) {
                log.warn("Duplicate event: {}", event.getEventId());
                return;
            }
            
            // 处理库存扣减
            inventoryService.deductInventory(event.getOrderId());
            
            // 记录去重
            eventDedupService.markProcessed(event.getEventId());
        };
    }
}
# application.yml - Stream配置
spring:
  cloud:
    function:
      definition: processOrderCreated;processPaymentCompleted
    stream:
      kafka:
        binder:
          brokers: broker1:9092,broker2:9092
          configuration:
            enable.idempotence: true
            acks: all
      bindings:
        orderEvents-out-0:
          destination: order.events.v1
          producer:
            useNativeEncoding: true
        processOrderCreated-in-0:
          destination: order.events.v1
          group: inventory-service
          consumer:
            enableDlq: true
            dlqName: order.events.v1.dlq

四、分布式事务的Saga模式

4.1 Choreography(编舞)vs Orchestration(编排)

【Choreography 编舞模式——各有各的舞步】

  订单服务                    库存服务                    支付服务
     │                          │                          │
     │── order.created ────────►│                          │
     │                          │── 扣库存                  │
     │                          │── inventory.reduced ────►│
     │                          │                          │── 发起支付
     │                          │                          │── payment.success ──►
     ◄── 收到通知 ──────────────┼──────────────────────────┘
     │── order.paid             │                          │
     │                          │                          │
  优点: 简单、无中心节点、天然解耦
  缺点: 流程散落在各服务中,难以追踪全局状态


【Orchestration 编排模式——一个指挥官统一调度】

                    ┌──────────────┐
                    │ Saga编排服务  │
                    │ (Saga        │
                    │  Orchestrator)│
                    └──────┬───────┘
                           │
              ┌────────────┼────────────┐
              │ 命令        │ 命令        │ 命令
              ▼            ▼            ▼
        ┌──────────┐ ┌──────────┐ ┌──────────┐
        │ 订单服务  │ │ 库存服务  │ │ 支付服务  │
        └──────────┘ └──────────┘ └──────────┘
              │            │            │
              │ 响应        │ 响应        │ 响应
              └────────────┼────────────┘
                           ▼
                    ┌──────────────┐
                    │ Saga编排服务  │
                    │ 决定下一步/   │
                    │ 补偿回滚      │
                    └──────────────┘

  优点: 流程集中管理、补偿逻辑清晰、易于追踪
  缺点: 编排服务成为关键路径、可能成为瓶颈

4.2 Saga模式对比

维度Choreography(编舞)Orchestration(编排)
复杂度低(每个服务有自己的舞步)中(需要定义状态机)
可追踪性差(事件散落在各Topic)好(编排服务集中记录)
耦合度极低低(只依赖编排服务)
补偿回滚难(补偿逻辑分散)易(编排服务统一调度)
单点故障编排服务是单点
适用规模小到中型流程(3-5步)中到大型流程(5+步)
推荐场景简单流程、团队规模小复杂流程、需要全局视图

4.3 Choreography模式的补偿实现

// 每个服务自己处理补偿
@Component
public class InventorySagaHandler {
    
    @KafkaListener(topics = "order.events.v1", groupId = "inventory-saga")
    public void handleOrderCreated(OrderCreatedEvent event) {
        try {
            // 扣减库存
            inventoryService.deduct(event.getOrderId(), event.getItems());
            
            // 成功 → 发布成功事件
            sagaEventPublisher.publish(new InventoryReducedEvent(event));
        } catch (InsufficientStockException e) {
            // 失败 → 发布失败事件(触发其他服务的补偿)
            sagaEventPublisher.publish(new InventoryReduceFailedEvent(event, e));
        }
    }
    
    @KafkaListener(topics = "order.events.v1", groupId = "inventory-saga")
    public void handleOrderCanceled(OrderCanceledEvent event) {
        // 补偿操作:恢复库存
        inventoryService.restore(event.getOrderId(), event.getItems());
    }
}

4.4 Orchestration模式的Saga状态机

@Component
public class OrderSagaOrchestrator {
    
    private final Map<String, SagaState> sagas = new ConcurrentHashMap<>();
    
    public enum SagaStep {
        VALIDATE_ORDER,      // 验证订单
        DEDUCT_INVENTORY,    // 扣库存
        PROCESS_PAYMENT,     // 处理支付
        NOTIFY_CUSTOMER,     // 通知用户
        COMPLETED,           // 完成
        COMPENSATING,        // 补偿中
        FAILED               // 失败
    }
    
    public void onOrderCreated(OrderCreatedEvent event) {
        SagaState state = new SagaState(event.getOrderId(), SagaStep.VALIDATE_ORDER);
        sagas.put(event.getOrderId(), state);
        
        // 进入下一状态:扣库存
        state.nextStep(SagaStep.DEDUCT_INVENTORY);
        sagaCommandPublisher.sendCommand(new DeductInventoryCommand(event));
    }
    
    public void onInventoryReduced(InventoryReducedEvent event) {
        SagaState state = sagas.get(event.getOrderId());
        if (state == null) return;
        
        // 扣库存成功 → 进入支付状态
        state.nextStep(SagaStep.PROCESS_PAYMENT);
        sagaCommandPublisher.sendCommand(new ProcessPaymentCommand(event));
    }
    
    public void onInventoryReduceFailed(InventoryReduceFailedEvent event) {
        SagaState state = sagas.get(event.getOrderId());
        if (state == null) return;
        
        // 扣库存失败 → 开始补偿
        state.nextStep(SagaStep.COMPENSATING);
        sagaCommandPublisher.sendCommand(new CancelOrderCommand(event.getOrderId()));
    }
}

五、事件溯源(Event Sourcing)

5.1 什么是事件溯源

传统的CRUD只存"当前状态",Event Sourcing存的是"所有变更事件"。

【CRUD vs Event Sourcing】

CRUD(存当前状态):
  ┌──────────────────────┐
  │  orders 表             │
  │  id=1, status=PAID,   │
  │  amount=299.00        │
  └──────────────────────┘
  只知道"现在是什么状态",不知道"怎么变成这样的"


Event Sourcing(存事件序列):
  ┌──────────────────────────────────────────┐
  │  order-events Topic                       │
  │                                            │
  │  ┌────────────────────────────────────┐   │
  │  │ eventId=1: order.created          │   │
  │  │ eventId=2: order.paid             │   │
  │  │ eventId=3: order.shipped          │   │
  │  │ eventId=4: order.delivered        │   │
  │  └────────────────────────────────────┘   │
  │                                            │
  │  重放事件序列就能得到任意时刻的状态!       │
  │  → 天然审计日志                            │
  │  → 任意时间点快照                          │
  │  → Bug修复后可以重放纠正                    │
  └──────────────────────────────────────────┘

5.2 事件溯源实现

// 订单聚合根(Event Sourcing 风格)
public class OrderAggregate {
    private String orderId;
    private OrderStatus status;
    private long amount;
    private List<OrderEvent> uncommittedEvents = new ArrayList<>();
    
    // 通过事件重放恢复状态(Snapshot + 事件重放)
    public static OrderAggregate fromHistory(List<OrderEvent> events) {
        OrderAggregate order = new OrderAggregate();
        for (OrderEvent event : events) {
            order.apply(event);  // 重放每个事件
        }
        return order;
    }
    
    // 创建订单
    public void create(String orderId, String userId, long amount) {
        OrderCreatedEvent event = new OrderCreatedEvent(orderId, userId, amount);
        apply(event);               // 应用事件,改变内存状态
        uncommittedEvents.add(event);  // 记录未提交事件
    }
    
    // 支付订单
    public void pay() {
        if (status != OrderStatus.CREATED) {
            throw new IllegalStateException("订单状态不允许支付");
        }
        OrderPaidEvent event = new OrderPaidEvent(orderId);
        apply(event);
        uncommittedEvents.add(event);
    }
    
    // 应用事件(内部状态转换)
    private void apply(OrderEvent event) {
        if (event instanceof OrderCreatedEvent e) {
            this.orderId = e.getOrderId();
            this.status = OrderStatus.CREATED;
            this.amount = e.getAmount();
        } else if (event instanceof OrderPaidEvent e) {
            this.status = OrderStatus.PAID;
        }
        // ... 其他事件处理
    }
    
    // 获取未提交事件
    public List<OrderEvent> getUncommittedEvents() {
        return Collections.unmodifiableList(uncommittedEvents);
    }
}

// 事件存储(写入Kafka)
@Service
public class OrderEventStore {
    
    private final KafkaTemplate<String, OrderEvent> kafkaTemplate;
    
    @Transactional
    public void saveEvents(OrderAggregate order) {
        for (OrderEvent event : order.getUncommittedEvents()) {
            kafkaTemplate.send("order.events.v1", 
                event.getOrderId(),  // 同一订单的事件到同一分区 → 保证顺序
                event
            ).get();  // 同步等待,确保写入成功
        }
    }
}

// 读取状态(快照 + 增量重放)
@Service
public class OrderRepository {
    
    @KafkaListener(topics = "order.events.v1", groupId = "order-read-model")
    public void updateReadModel(OrderEvent event) {
        // 维护一个读模型(比如写入MySQL或Redis)
        // 每次收到事件就更新读模型
        jdbcTemplate.update("""
            INSERT INTO order_view (order_id, status, amount, version)
            VALUES (?, ?, ?, 1)
            ON DUPLICATE KEY UPDATE 
                status = ?, 
                amount = ?,
                version = version + 1
            """, 
            event.getOrderId(), event.getOrderStatus(), 
            event.getAmount(), /* ... */
        );
    }
}

5.3 CQRS与事件溯源

【CQRS + Event Sourcing 架构图】

              ┌──────────────────────────────┐
              │       Kafka 事件总线           │
              │                                │
  命令模型     │  ┌──────────────────────────┐ │    查询模型
  (写)        │  │  order.events.v1          │ │    (读)
              │  │  ┌──────────────────────┐ │ │
  ┌────────┐  │  │  │ event1: created     │ │ │  ┌────────┐
  │Order   │──┼──│  │ event2: paid        │ │ │  │MySQL   │
  │Command │  │  │  │ event3: shipped     │─┼─┼──►│读模型   │──►前端展示
  │Handler │  │  │  └──────────────────────┘ │ │  │order_view│
  └────────┘  │  └──────────────────────────┘ │  └────────┘
              │                                │
              │  ┌──────────────────────────┐ │  ┌────────┐
              │  │  inventory.events.v1     │─┼──►│ES      │──►搜索
              │  │  payment.events.v1       │─┼──►│Redis   │──►缓存
              │  └──────────────────────────┘ │  └────────┘
              └──────────────────────────────┘

  核心思想:
  • 命令端(写)只关心业务逻辑和事件生成
  • 查询端(读)可以根据不同场景构建不同的读模型
  • 事件是唯一的数据源(Single Source of Truth)

六、事件驱动架构的常见反模式

反模式症状正确做法
一切皆事件每个小操作都发事件,Topic爆炸只在业务关键节点发事件
空事件事件只含ID,消费者要从数据库查数据事件应自描述,包含处理所需的关键数据
同步等待事件发了事件后阻塞等待响应事件驱动是异步的,用Saga或回调处理
过度自信幂等假定消息不会重复而不做去重永远假设消息可能重复
事件中放敏感数据订单事件包含完整地址、手机号仅放必要数据,敏感字段需脱敏
忽略事件顺序同一实体的多个事件散到不同分区用实体ID做分区Key保证有序

本篇小结

事件驱动架构不是银弹,但在正确的场景下它能把微服务的复杂度降一个数量级:

  • 事件驱动 vs 请求响应:当你需要3个以上下游服务协作时,事件驱动是更好的选择——发布者不关心消费者,天然解耦
  • 领域事件设计:用过去时命名,自描述但不臃肿,带eventId用于去重,通过version字段支持格式演进
  • Kafka事件总线:每类业务一个Topic,用实体ID做分区Key保证有序,Kafka的分区天然支持消费者的水平扩展
  • Saga分布式事务:3-5步简单流程用Choreography(编舞),5步以上复杂流程用Orchestration(编排),补偿逻辑是Saga的核心
  • 事件溯源:存储事件序列而非当前状态,天然审计+任意时间点恢复+CQRS读写分离——但运维复杂度也更高

90篇文章,从Kafka是什么到事件驱动微服务架构,这个系列到这里就结束了。希望这些内容能帮你从"会用Kafka"进化到"能用好Kafka"。技术之路无穷尽,祝你越走越远。


上一篇【第89篇】实时数据同步平台的Kafka实战——MySQL CDC与Kafka的最佳组合
系列完结,感谢阅读


更多推荐