事件驱动架构 EDA 深度实战:Saga、事件溯源与 CQRS

事件驱动架构完整落地:事件总线设计、Saga 分布式事务模式、事件溯源与 CQRS 读写分离,以及 DDD 聚合域事件与 Kafka 事件流编排。

事件驱动架构(Event-Driven Architecture, EDA)是现代分布式系统的核心通信范式。本文从事件语义拆解出发,逐步深入事件总线、Saga 分布式事务、事件溯源、CQRS、DDD 聚合事件与 Kafka Streams 流编排,提供可直接落地的代码、设计图与生产级实践。


一、EDA 核心语义:事件、消息与命令

在事件驱动系统中,EventMessageCommand 三个概念经常被混用。准确区分它们是设计清晰架构的第一步。

┌──────────────────────────────────────────────────────────────┐
│                      EDA 语义分层                             │
├──────────────────────────────────────────────────────────────┤
│  Command 命令     │  意图(Intent)  │  要求系统执行某个操作    │
│  Event  事件      │  事实(Fact)    │  某事已发生,不可变       │
│  Message 消息     │  载体(Carrier) │  事件/命令的传输信封      │
└──────────────────────────────────────────────────────────────┘

**命令(Command)**发送到特定目标,期望产生副作用;**事件(Event)**由系统发布,表示状态变更已发生;**消息(Message)**是二者在传输层上的统称。一个典型的交互流如下:

用户 ──[PlaceOrder Command]──> OrderService
OrderService ──[OrderCreated Event]──> EventBus
EventBus ──[Event Message]──> InventoryService / PaymentService / NotificationService

关键设计原则:

  • 命令是定向的(Addressed),可以失败、可以被拒绝。
  • 事件是广播的(Broadcasted),已被发布后不可撤回,消费端通过补偿处理异常。
  • 消息保证传输(At-Least-Once),不保证消费顺序(除非显式配置)。

二、事件总线设计:Pub-Sub 与队列语义

事件总线(Event Bus)是 EDA 的神经系统,决定事件的流转路径与消费语义。核心有两种模型:

维度发布-订阅(Pub-Sub)队列(Queue)
投递模型广播给所有订阅者竞争消费,单条消息仅被一个消费者处理
典型场景订单创建后同时通知库存、支付、物流异步任务处理,如生成报表、发送邮件
背压控制依赖消费者自身速率可通过队列深度 + 消费者数调控
代表中间件Kafka、Redis PubSub、RabbitMQ FanoutRabbitMQ Queue、RocketMQ、AWS SQS
消费偏移每个订阅者独立维护 offset队列消费后删除或归档

推荐混合架构: 使用 Kafka 的 Topic-Partition 机制同时实现广播与分区消费:不同 Consumer Group 订阅同一 Topic 实现 Pub-Sub;同一 Group 内的消费者竞争消费同一 Partition 实现 Queue 语义。

2.1 轻量级事件总线实现(Java)

package com.example.eda.bus;

import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.function.Consumer;

/**
 * 内存级事件总线,用于单元测试或单机事件编排演示
 * 生产环境应替换为 Kafka / RabbitMQ / Pulsar
 */
public class InMemoryEventBus implements EventBus {

    // 按事件类型存储监听器,线程安全
    private final Map<Class<?>, CopyOnWriteArrayList<Consumer<Object>>> subscribers
            = new ConcurrentHashMap<>();

    // 发布事件:广播给所有订阅了该事件类型的消费者
    @Override
    @SuppressWarnings("unchecked")
    public <T> void publish(T event) {
        if (event == null) return;
        Class<?> eventType = event.getClass();
        // 同时匹配精确类型与父类监听器
        subscribers.forEach((type, listeners) -> {
            if (type.isAssignableFrom(eventType)) {
                listeners.forEach(listener -> {
                    // 异步投递避免阻塞发布者
                    EventBusExecutors.submit(() -> listener.accept(event));
                });
            }
        });
    }

    // 订阅事件
    @Override
    @SuppressWarnings("unchecked")
    public <T> void subscribe(Class<T> eventType, Consumer<T> listener) {
        subscribers.computeIfAbsent(eventType, k -> new CopyOnWriteArrayList<>())
                   .add((Consumer<Object>) listener);
    }
}

2.2 事件信封封装

package com.example.eda.bus;

import java.time.Instant;
import java.util.UUID;

/**
 * 事件信封:携带元数据,用于链路追踪与幂等判断
 */
public record EventEnvelope<T>(
    String eventId,           // 全局唯一事件 ID
    String correlationId,     // 链路追踪 ID
    String sourceService,     // 产生事件的服务名
    Instant occurredOn,       // 事件发生时间(业务时间)
    T payload                 // 实际负载
) {
    public static <T> EventEnvelope<T> wrap(String source, String correlationId, T payload) {
        return new EventEnvelope<>(
            UUID.randomUUID().toString(),
            correlationId != null ? correlationId : UUID.randomUUID().toString(),
            source,
            Instant.now(),
            payload
        );
    }
}

三、Saga 分布式事务模式

在微服务中,ACID 事务跨服务不可用。Saga 模式将长事务拆分为本地事务序列,每个本地事务提交后立即发布事件驱动下一步;失败时通过补偿事务回滚已完成的操作。

Saga 有两种实现风格:

3.1 编排式 Saga(Choreography)

每个服务完成本地事务后发送事件,下游服务监听事件自主决策。无中央协调器,松耦合但流程分散。

package com.example.eda.saga;

import com.example.eda.bus.EventBus;
import com.example.eda.bus.EventEnvelope;

/**
 * 编排式 Saga:库存服务监听订单创建事件,扣减库存后发布 InventoryReserved
 */
public class InventoryChoreographyHandler {

    private final EventBus eventBus;
    private final InventoryRepository inventoryRepository;

    public InventoryChoreographyHandler(EventBus eventBus, InventoryRepository repo) {
        this.eventBus = eventBus;
        this.inventoryRepository = repo;
    }

    // 订阅订单创建事件
    public void onOrderCreated(EventEnvelope<OrderCreatedEvent> envelope) {
        var event = envelope.payload();
        String orderId = event.orderId();
        String sku = event.sku();
        int qty = event.quantity();

        try {
            // 本地事务:扣减库存
            inventoryRepository.decrease(sku, qty);

            // 发布库存预留成功事件,支付服务将监听此事件
            var reserved = new InventoryReservedEvent(orderId, sku, qty);
            eventBus.publish(EventEnvelope.wrap("inventory", envelope.correlationId(), reserved));

        } catch (InsufficientStockException e) {
            // 发布补偿事件,触发订单取消
            var failed = new InventoryReservationFailedEvent(orderId, sku, qty, e.getMessage());
            eventBus.publish(EventEnvelope.wrap("inventory", envelope.correlationId(), failed));
        }
    }
}

3.2 编排式 Saga:支付服务补偿示例

package com.example.eda.saga;

/**
 * 支付服务监听库存预留成功事件,完成扣款;若后续步骤失败需触发退款补偿
 */
public class PaymentChoreographyHandler {

    private final EventBus eventBus;
    private final PaymentGateway paymentGateway;

    public void onInventoryReserved(EventEnvelope<InventoryReservedEvent> envelope) {
        var event = envelope.payload();
        String orderId = event.orderId();

        try {
            String transactionId = paymentGateway.charge(orderId, event.amount());
            var paid = new PaymentCompletedEvent(orderId, transactionId);
            eventBus.publish(EventEnvelope.wrap("payment", envelope.correlationId(), paid));
        } catch (PaymentException e) {
            // 支付失败,发布事件通知库存服务释放库存(补偿)
            var failed = new PaymentFailedEvent(orderId, e.getMessage());
            eventBus.publish(EventEnvelope.wrap("payment", envelope.correlationId(), failed));
        }
    }

    // 监听 Saga 失败信号,执行退款补偿
    public void onSagaFailed(EventEnvelope<SagaFailedEvent> envelope) {
        var event = envelope.payload();
        // 幂等退款:通过 transactionId 避免重复退款
        paymentGateway.refund(event.transactionId());
    }
}

3.3 协调式 Saga(Orchestration)

由中央 Saga 协调器统一下发命令到各服务,维护状态机。流程集中可控,适合复杂长事务。

package com.example.eda.saga

import kotlinx.coroutines.*

/**
 * 协调式 Saga:OrderSagaOrchestrator 作为中央状态机驱动各服务
 * Kotlin + 挂起函数实现异步非阻塞编排
 */
class OrderSagaOrchestrator(
    private val inventoryService: InventoryService,
    private val paymentService: PaymentService,
    private val shippingService: ShippingService,
    private val eventBus: EventBus
) {

    // 发起 Saga
    suspend fun startSaga(orderId: String, sku: String, qty: Int, amount: BigDecimal): SagaResult {
        val correlationId = generateCorrelationId()
        val log = mutableListOf<SagaStep>()

        return try {
            // 第一步:预留库存
            val reserved = inventoryService.reserve(orderId, sku, qty)
            log.add(SagaStep("INVENTORY_RESERVED", reserved))

            // 第二步:扣款
            val paid = paymentService.charge(orderId, amount)
            log.add(SagaStep("PAYMENT_COMPLETED", paid))

            // 第三步:创建物流单
            val shipped = shippingService.createShipment(orderId, sku, qty)
            log.add(SagaStep("SHIPMENT_CREATED", shipped))

            SagaResult.Success(correlationId)

        } catch (e: Exception) {
            // 反向补偿:按相反顺序执行补偿操作
            compensate(log, orderId)
            SagaResult.Failed(correlationId, e.message)
        }
    }

    // 补偿逻辑:逆序回滚已完成的步骤
    private suspend fun compensate(log: List<SagaStep>, orderId: String) {
        log.asReversed().forEach { step ->
            when (step.name) {
                "SHIPMENT_CREATED" -> shippingService.cancelShipment(orderId)
                "PAYMENT_COMPLETED" -> paymentService.refund(step.result.transactionId)
                "INVENTORY_RESERVED" -> inventoryService.release(orderId)
            }
        }
    }

    private fun generateCorrelationId(): String = java.util.UUID.randomUUID().toString()
}

// Saga 执行结果密封类
data class SagaResult(val correlationId: String) {
    class Success(correlationId: String) : SagaResult(correlationId)
    class Failed(correlationId: String, val reason: String?) : SagaResult(correlationId)
}

四、事件溯源(Event Sourcing)

事件溯源将系统状态存储为一系列不可变事件,而非直接保存当前状态。通过重放事件可还原任意时刻的系统状态。

4.1 核心概念

  • Event Store:仅追加的事件存储,支持按聚合 ID 顺序读取。
  • Aggregate:业务一致性边界,通过重放事件还原自身状态。
  • Snapshot:聚合事件过多时保存的状态快照,加速重放。
  • Projection:读模型,通过监听事件构建,支持任意查询维度。

4.2 事件存储实现

package com.example.eda.eventsourcing;

import java.util.List;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.stream.Collectors;

/**
 * 内存事件存储演示,生产环境使用 EventStoreDB / Apache Pulsar / PostgreSQL 仅追加表
 */
public class InMemoryEventStore implements EventStore {

    // 按聚合 ID 分组存储有序事件列表
    private final Map<String, List<DomainEvent>> streams = new ConcurrentHashMap<>();
    private final AtomicLong globalVersion = new AtomicLong(0);

    // 持久化事件,保证聚合内顺序写入
    @Override
    public synchronized void append(String aggregateId, List<DomainEvent> events, long expectedVersion) {
        List<DomainEvent> stream = streams.computeIfAbsent(aggregateId, k -> new CopyOnWriteArrayList<>());

        // 乐观并发控制:版本号冲突时拒绝写入
        if (stream.size() != expectedVersion) {
            throw new ConcurrencyException(
                "聚合 " + aggregateId + " 版本冲突,期望 " + expectedVersion + ",实际 " + stream.size()
            );
        }

        for (DomainEvent event : events) {
            // 为每个事件分配全局版本号与聚合内版本号
            long version = globalVersion.incrementAndGet();
            EventMeta meta = new EventMeta(version, aggregateId, stream.size() + 1, Instant.now());
            event.assignMeta(meta);
            stream.add(event);
        }
    }

    // 读取聚合的全部事件,用于状态重放
    @Override
    public List<DomainEvent> readStream(String aggregateId) {
        return streams.getOrDefault(aggregateId, List.of())
                      .stream()
                      .map(e -> (DomainEvent) e) // 类型恢复
                      .collect(Collectors.toList());
    }

    // 从指定版本号开始读取(配合快照使用)
    @Override
    public List<DomainEvent> readStreamFrom(String aggregateId, long fromVersion) {
        List<DomainEvent> stream = streams.getOrDefault(aggregateId, List.of());
        return stream.stream()
                     .filter(e -> e.meta().aggregateVersion() >= fromVersion)
                     .collect(Collectors.toList());
    }
}

4.3 聚合根与事件重放

package com.example.eda.eventsourcing;

/**
 * 订单聚合根:通过重放事件还原状态,命令执行后产生新事件
 */
public class OrderAggregate {

    private String orderId;
    private OrderStatus status;
    private List<OrderLine> lines = new ArrayList<>();
    private long version; // 当前聚合版本号,用于乐观并发控制

    // 空构造用于反射重建
    public OrderAggregate() {}

    // 通过事件流还原聚合状态
    public void rehydrate(List<DomainEvent> events) {
        for (DomainEvent event : events) {
            apply(event); // 无分支,仅状态变更
            this.version = event.meta().aggregateVersion();
        }
    }

    // 处理命令:创建订单
    public List<DomainEvent> create(String orderId, List<OrderLine> lines) {
        if (this.orderId != null) {
            throw new IllegalStateException("订单已存在");
        }
        return List.of(new OrderCreatedEvent(orderId, lines, Instant.now()));
    }

    // 处理命令:确认支付
    public List<DomainEvent> confirmPayment(String transactionId) {
        if (status != OrderStatus.PENDING_PAYMENT) {
            throw new IllegalStateException("订单状态不支持确认支付: " + status);
        }
        return List.of(new PaymentConfirmedEvent(orderId, transactionId, Instant.now()));
    }

    // 事件路由:根据事件类型调用对应 apply 方法
    private void apply(DomainEvent event) {
        switch (event) {
            case OrderCreatedEvent e -> {
                this.orderId = e.orderId();
                this.lines = new ArrayList<>(e.lines());
                this.status = OrderStatus.PENDING_PAYMENT;
            }
            case PaymentConfirmedEvent e -> {
                this.status = OrderStatus.PAID;
            }
            case OrderShippedEvent e -> {
                this.status = OrderStatus.SHIPPED;
            }
            default -> throw new IllegalArgumentException("未知事件类型: " + event.getClass());
        }
    }

    public long version() { return version; }
    public OrderStatus status() { return status; }
}

4.4 快照机制

package com.example.eda.eventsourcing;

/**
 * 快照存储:当聚合事件数超过阈值时保存当前状态,避免全量重放
 */
public class SnapshotRepository {

    private final Map<String, Snapshot> snapshots = new ConcurrentHashMap<>();
    private static final int SNAPSHOT_THRESHOLD = 50; // 每 50 个事件打一次快照

    public void save(OrderAggregate aggregate) {
        // 序列化当前状态为快照
        Snapshot snapshot = new Snapshot(
            aggregate.getOrderId(),
            aggregate.version(),
            serialize(aggregate)
        );
        snapshots.put(aggregate.getOrderId(), snapshot);
    }

    // 恢复快照后仅需重放后续事件
    public Optional<OrderAggregate> load(String aggregateId, EventStore eventStore) {
        Snapshot snapshot = snapshots.get(aggregateId);
        OrderAggregate aggregate = new OrderAggregate();

        List<DomainEvent> events;
        if (snapshot != null) {
            aggregate = deserialize(snapshot.data());
            // 从快照版本之后继续读取
            events = eventStore.readStreamFrom(aggregateId, snapshot.version() + 1);
        } else {
            events = eventStore.readStream(aggregateId);
        }

        if (events.isEmpty() && snapshot == null) {
            return Optional.empty();
        }

        aggregate.rehydrate(events);
        return Optional.of(aggregate);
    }

    // 判断是否需要打快照
    public boolean shouldSnapshot(long currentVersion) {
        return currentVersion > 0 && currentVersion % SNAPSHOT_THRESHOLD == 0;
    }
}

五、CQRS:命令查询职责分离

CQRS 将读模型与写模型分离,写模型专注于业务不变式与事件生成,读模型通过投影构建为查询高度优化的视图。

┌──────────────┐     Command      ┌─────────────┐     Event      ┌────────────────┐
│   客户端      │ ───────────────> │  Command    │ ────────────> │   Event Store  │
│  (写请求)     │                  │  Handler    │               │  (唯一真相源)   │
└──────────────┘                  └─────────────┘               └────────────────┘
                                                                       │
                                                         ┌─────────────┼─────────────┐
                                                         ▼             ▼             ▼
                                                   ┌─────────┐   ┌─────────┐   ┌─────────┐
                                                   │Projection│   │Projection│   │Projection│
                                                   │ 订单列表  │   │ 订单统计  │   │ 用户视图  │
                                                   └─────────┘   └─────────┘   └─────────┘
┌──────────────┐     Query        ┌─────────────┐         ▲           ▲           ▲
│   客户端      │ ───────────────> │  Query      │─────────┘───────────┘───────────┘
│  (读请求)     │                  │  Handler    │      (Read Model: 任意存储)
└──────────────┘                  └─────────────┘

5.1 写模型侧:命令处理器

package com.example.eda.cqrs;

/**
 * 写侧命令处理器:负责校验、执行业务逻辑、持久化事件
 */
public class OrderCommandHandler {

    private final EventStore eventStore;
    private final SnapshotRepository snapshotRepository;

    public void handle(CreateOrderCommand cmd) {
        // 从事件存储还原聚合(或新建)
        OrderAggregate aggregate = snapshotRepository
            .load(cmd.orderId(), eventStore)
            .orElseGet(OrderAggregate::new);

        // 执行业务命令,生成事件(不写 DB,只生成事件)
        List<DomainEvent> events = aggregate.create(cmd.orderId(), cmd.lines());

        // 将事件追加到事件存储
        eventStore.append(cmd.orderId(), events, aggregate.version());

        // 检查是否需要打快照
        long newVersion = aggregate.version() + events.size();
        if (snapshotRepository.shouldSnapshot(newVersion)) {
            // 重放后保存快照
            OrderAggregate fresh = snapshotRepository.load(cmd.orderId(), eventStore).orElseThrow();
            snapshotRepository.save(fresh);
        }
    }
}

5.2 读模型侧:投影处理器

package com.example.eda.cqrs;

import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;

/**
 * 读侧投影:监听事件,构建面向查询的扁平化视图
 * 可存储于 Elasticsearch / Redis / MongoDB / 关系型数据库读库
 */
public class OrderListProjection implements EventHandler {

    // 模拟读模型存储,生产环境使用独立数据库
    private final Map<String, OrderListView> readModel = new ConcurrentHashMap<>();

    @Override
    public void on(OrderCreatedEvent event) {
        OrderListView view = new OrderListView(
            event.orderId(),
            event.lines().stream().mapToInt(OrderLine::quantity).sum(),
            event.lines().stream().map(OrderLine::price).reduce(BigDecimal.ZERO, BigDecimal::add),
            "PENDING_PAYMENT",
            event.occurredOn()
        );
        readModel.put(event.orderId(), view);
        indexToElasticsearch(view); // 同步或异步索引到搜索引擎
    }

    @Override
    public void on(PaymentConfirmedEvent event) {
        OrderListView view = readModel.get(event.orderId());
        if (view != null) {
            view = view.withStatus("PAID").withPaidAt(event.occurredOn());
            readModel.put(event.orderId(), view);
            updateElasticsearch(view);
        }
    }

    @Override
    public void on(OrderShippedEvent event) {
        OrderListView view = readModel.get(event.orderId());
        if (view != null) {
            view = view.withStatus("SHIPPED");
            readModel.put(event.orderId(), view);
            updateElasticsearch(view);
        }
    }

    // 查询接口:支持分页、过滤、排序(由读库能力决定)
    public List<OrderListView> query(OrderQueryCriteria criteria) {
        return readModel.values().stream()
            .filter(criteria::matches)
            .sorted(criteria.getComparator())
            .skip(criteria.offset())
            .limit(criteria.limit())
            .collect(Collectors.toList());
    }
}

六、DDD 聚合与领域事件

DDD(领域驱动设计)中的聚合(Aggregate)是事件溯源与 CQRS 的完美宿主。领域事件从聚合内部产生,代表业务上真正有意义的状态变更。

6.1 领域事件定义

package com.example.eda.domain;

/**
 * 领域事件基类:所有业务事件继承此类
 * 事件命名应采用过去式,表达"已发生"的不可变事实
 */
public abstract class DomainEvent {

    private EventMeta meta; // 元数据在持久化时注入

    public void assignMeta(EventMeta meta) {
        if (this.meta != null) throw new IllegalStateException("事件元数据只能赋值一次");
        this.meta = meta;
    }

    public EventMeta meta() { return meta; }

    // 业务发生时间应由聚合在创建事件时指定,而非系统自动注入
    public abstract Instant occurredOn();
}

// 具体领域事件
public record OrderCreatedEvent(
    String orderId,
    List<OrderLine> lines,
    Instant occurredOn
) extends DomainEvent {}

public record PaymentConfirmedEvent(
    String orderId,
    String transactionId,
    Instant occurredOn
) extends DomainEvent {}

6.2 聚合边界与事务一致性

package com.example.eda.domain;

/**
 * 账户聚合:演示聚合内强一致性,聚合间最终一致性
 */
public class BankAccountAggregate {

    private String accountId;
    private BigDecimal balance;
    private List<PendingTransfer> pendingTransfers = new ArrayList<>();

    // 同一聚合内转账:命令执行后立即产生事件,强一致
    public List<DomainEvent> transferWithinLimit(String toAccount, BigDecimal amount, String transferId) {
        if (balance.compareTo(amount) < 0) {
            throw new InsufficientBalanceException("余额不足");
        }
        if (amount.compareTo(new BigDecimal("100000")) > 0) {
            throw new LimitExceededException("单笔转账超过限额");
        }
        return List.of(
            new TransferInitiatedEvent(accountId, toAccount, amount, transferId, Instant.now()),
            new BalanceDebitedEvent(accountId, amount, balance.subtract(amount), Instant.now())
        );
    }

    // 跨聚合转账:仅在本聚合生成事件,接收方通过事件监听异步处理,最终一致
    public List<DomainEvent> initiateCrossAggregateTransfer(String toAccount, BigDecimal amount, String transferId) {
        // 校验与事件生成同上
        return transferWithinLimit(toAccount, amount, transferId);
        // 接收方聚合将在独立事务中处理 TransferReceivedEvent
    }
}

6.3 应用服务层:协调聚合与基础设施

package com.example.eda.application;

/**
 * 应用服务:无业务逻辑,仅负责编排领域层与基础设施层
 */
@Service
public class OrderApplicationService {

    private final EventStore eventStore;
    private final SnapshotRepository snapshotRepository;
    private final DomainEventPublisher publisher;

    @Transactional // 写数据库事务(非分布式事务)
    public void createOrder(CreateOrderCommand cmd) {
        // 1. 还原聚合
        OrderAggregate order = snapshotRepository
            .load(cmd.orderId(), eventStore)
            .orElseGet(OrderAggregate::new);

        // 2. 执行领域逻辑
        List<DomainEvent> events = order.create(cmd.orderId(), cmd.lines());

        // 3. 持久化事件(写事件存储)
        eventStore.append(cmd.orderId(), events, order.version());

        // 4. 发布事件到外部队列,供投影与其他服务消费
        events.forEach(publisher::publish);
    }
}

七、Kafka Streams 事件流编排

Apache Kafka 不仅是消息队列,其 Streams API 提供了轻量级流处理能力,可在事件管道中实现状态化计算、窗口聚合与多流 Join。

7.1 流拓扑定义:订单金额实时统计

package com.example.eda.streams;

import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.KafkaStreams;
import org.apache.kafka.streams.StreamsBuilder;
import org.apache.kafka.streams.kstream.*;
import org.apache.kafka.streams.state.Stores;

import java.time.Duration;

/**
 * Kafka Streams 拓扑:实时统计每 5 分钟窗口内的订单总金额与订单数
 */
public class OrderAnalyticsTopology {

    public Topology build() {
        StreamsBuilder builder = new StreamsBuilder();

        // 定义事件流:从 order-events Topic 读取
        KStream<String, OrderEvent> orderStream = builder.stream(
            "order-events",
            Consumed.with(Serdes.String(), new OrderEventSerde())
        );

        // 过滤出已支付事件,按商品类目分组,滑动窗口聚合
        orderStream
            .filter((key, event) -> event instanceof PaymentConfirmedEvent)
            .groupBy((key, event) -> event.category(), Grouped.with(Serdes.String(), new OrderEventSerde()))
            .windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(5)))
            .aggregate(
                OrderStats::new, // 初始值
                (category, event, stats) -> stats.accumulate(event.amount()), // 累加
                Materialized.<String, OrderStats>as(Stores.inMemoryWindowStore(
                    "order-stats-store",
                    Duration.ofHours(24), // 保留窗口 24 小时
                    Duration.ofMinutes(5),
                    false
                ))
                .withKeySerde(Serdes.String())
                .withValueSerde(new OrderStatsSerde())
            )
            .toStream()
            .to("order-analytics", Produced.with(new WindowedSerde<>(Serdes.String()), new OrderStatsSerde()));

        // 多流 Join:订单流与物流流通过订单 ID 关联,计算端到端履约时长
        KTable<String, ShipmentEvent> shipmentTable = builder.table(
            "shipment-events",
            Consumed.with(Serdes.String(), new ShipmentEventSerde())
        );

        orderStream
            .filter((k, e) -> e instanceof OrderCreatedEvent)
            .selectKey((k, e) -> e.orderId())
            .join(
                shipmentTable,
                (order, shipment) -> new FulfillmentMetrics(
                    order.orderId(),
                    Duration.between(order.occurredOn(), shipment.occurredOn()).toMinutes()
                ),
                Joined.with(Serdes.String(), new OrderEventSerde(), new ShipmentEventSerde())
            )
            .to("fulfillment-metrics", Produced.with(Serdes.String(), new MetricsSerde()));

        return builder.build();
    }
}

7.2 Kafka Streams 配置与启动

package com.example.eda.streams;

import org.apache.kafka.streams.StreamsConfig;
import java.util.Properties;

public class KafkaStreamsConfig {

    public static Properties configure(String appId, String bootstrapServers) {
        Properties props = new Properties();
        props.put(StreamsConfig.APPLICATION_ID_CONFIG, appId);
        props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
        props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
        props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, OrderEventSerde.class);
        props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, StreamsConfig.EXACTLY_ONCE_V2);
        props.put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 5000); // 5 秒提交一次
        props.put(StreamsConfig.NUM_STREAM_THREADS_CONFIG, 4);    // 并行线程数
        props.put(StreamsConfig.REPLICATION_FACTOR_CONFIG, 3);    // 状态存储副本数
        return props;
    }

    public static void main(String[] args) {
        Properties props = configure("order-analytics-app", "kafka:9092");
        Topology topology = new OrderAnalyticsTopology().build();
        KafkaStreams streams = new KafkaStreams(topology, props);

        Runtime.getRuntime().addShutdownHook(new Thread(streams::close));
        streams.start();
    }
}

八、EDA 与 RPC 对比分析

对比维度RPC / REST 同步调用事件驱动异步通信
耦合度服务间直接依赖,需知悉对方地址与接口契约通过事件总线解耦,生产者不知消费者存在
响应延迟低延迟,直接返回结果引入中间件延迟,适合容忍秒级延迟场景
可用性链路中任一服务故障导致级联失败中间件缓冲消息,消费者故障不影响生产者
事务一致性强一致(局部),跨服务需引入 Saga天然最终一致,通过 Saga 补偿保证业务完整
流量削峰需额外引入熔断限流队列天然削峰填谷,消费者按需伸缩
系统复杂度调用链路直观,调试简单事件传播路径隐式,需追踪与监控
数据一致性调试单链路断点即可需分布式追踪(TraceId)与事件回放
适用场景实时查询、强一致事务、短链路异步处理、长流程编排、高吞吐、多订阅者

架构选型建议: 读操作为主的查询优先使用 RPC / GraphQL;写操作涉及多服务协作或需异步处理的场景优先使用 EDA。混合架构下,API Gateway 接收外部请求,写路径转换为事件,读路径直接查询或访问 CQRS 读模型。


九、数据一致性策略

事件驱动系统的核心挑战是一致性:生产者写数据库与发事件之间、消费者处理事件与更新本地状态之间,均存在分布式事务难题。

9.1 发件箱模式(Outbox Pattern)

package com.example.eda.outbox;

import org.springframework.transaction.annotation.Transactional;

/**
 * 发件箱模式:业务表写入与事件记录在同一本地事务中,保证原子性
 */
@Service
public class OrderServiceWithOutbox {

    private final OrderRepository orderRepository;
    private final OutboxRepository outboxRepository;

    @Transactional
    public void placeOrder(PlaceOrderRequest request) {
        // 1. 写业务数据
        Order order = Order.create(request);
        orderRepository.save(order);

        // 2. 将待发送事件写入 outbox 表(同一数据库事务)
        OutboxRecord record = new OutboxRecord(
            UUID.randomUUID().toString(),
            "order-events",           // 目标 Topic
            order.getId(),            // 聚合 ID / Partition Key
            "OrderCreated",           // 事件类型
            toJson(new OrderCreatedEvent(order.getId(), order.getLines(), Instant.now())),
            OutboxStatus.PENDING,
            Instant.now()
        );
        outboxRepository.save(record);
        // 事务提交后,业务数据与事件记录必然同时存在
    }
}

9.2 发件箱轮询投递器

package com.example.eda.outbox;

import org.springframework.scheduling.annotation.Scheduled;

/**
 * 定时轮询 outbox 表,将未发送事件投递到 Kafka,成功后删除或标记为已发送
 */
@Component
public class OutboxPoller {

    private final OutboxRepository outboxRepository;
    private final KafkaTemplate<String, String> kafkaTemplate;

    // 每 500ms 轮询一次,生产环境可配合 Debezium 实现 CDC 实时捕获
    @Scheduled(fixedRate = 500)
    public void pollAndPublish() {
        List<OutboxRecord> pending = outboxRepository.findByStatus(OutboxStatus.PENDING, 100);

        for (OutboxRecord record : pending) {
            try {
                kafkaTemplate.send(record.topic(), record.aggregateId(), record.payload())
                    .whenComplete((result, ex) -> {
                        if (ex == null) {
                            outboxRepository.markAsSent(record.id());
                        } else {
                            outboxRepository.incrementRetry(record.id(), ex.getMessage());
                        }
                    });
            } catch (Exception e) {
                // 投递异常,留待下次轮询重试
                outboxRepository.incrementRetry(record.id(), e.getMessage());
            }
        }
    }
}

9.3 幂等消费

package com.example.eda.idempotency;

import org.springframework.stereotype.Component;
import java.util.concurrent.ConcurrentHashMap;

/**
 * 幂等键存储:基于事件 ID 去重,防止消费者重复处理导致状态不一致
 */
@Component
public class IdempotencyKeyStore {

    // 生产环境使用 Redis SETNX 或数据库唯一索引
    private final Set<String> processedIds = ConcurrentHashMap.newKeySet();

    public boolean isProcessed(String eventId) {
        return !processedIds.add(eventId); // 已存在返回 true(已处理)
    }
}

@Service
public class InventoryConsumer {

    private final IdempotencyKeyStore idempotency;
    private final InventoryRepository inventoryRepository;

    public void onInventoryReserved(EventEnvelope<InventoryReservedEvent> envelope) {
        String eventId = envelope.eventId();
        if (idempotency.isProcessed(eventId)) {
            // 已处理,直接返回,保证幂等
            return;
        }
        // 执行业务逻辑...
        inventoryRepository.decrease(envelope.payload().sku(), envelope.payload().qty());
    }
}

十、生产级实践

10.1 事件顺序性保障

  • 同一聚合内严格有序:使用 Kafka Partition Key = aggregateId,确保同一聚合事件进入同一 Partition,由单消费者顺序处理。
  • 跨聚合无序可接受:不同聚合之间天然无需全局顺序,可并行消费提升吞吐。
  • 因果顺序需求:对于跨聚合因果关系,下游服务通过 Saga 或等待依赖条件满足后再处理。
// Kafka Producer 按聚合 ID 分区,保证聚合内顺序
kafkaTemplate.send(
    ProducerRecord<>(
        "order-events",
        orderAggregateId,  // 分区键 = 聚合 ID
        eventPayload
    )
);

10.2 事件去重与至少一次语义

所有消费者必须实现幂等处理

  • 消费者侧存储已处理事件 ID(Redis / 数据库唯一键)。
  • 业务操作基于状态机判断,而非事件计数。例如库存扣减前检查订单状态是否为 PENDING
  • Kafka 配置 enable.idempotence=true,保证生产者端到端幂等。

10.3 事件版本演化(Versioning)

Schema 变更不可避免,推荐策略:

package com.example.eda.schema;

/**
 * 事件版本包装器:支持向后兼容的 schema 演化
 */
public record VersionedEvent(
    int schemaVersion,    // 当前 schema 版本号
    String eventType,     // 事件类型全限定名
    String payload        // JSON / Avro / Protobuf 序列化数据
) {
    public static final int CURRENT_VERSION = 2;
}
演化策略说明适用场景
新增可选字段旧消费者忽略未知字段,新消费者读取最常见的向后兼容变更
字段重命名保留旧字段别名,双写一段时间API 重构期过渡
事件类型升级发布 V2 事件类型,消费者逐步迁移语义发生本质变化
快照转换升级时读取旧事件,写入新 schema 快照大规模 schema 升级

10.4 可观测性建设

package com.example.eda.observability;

import io.micrometer.core.instrument.MeterRegistry;
import org.aspectj.lang.annotation.Aspect;

/**
 * 事件处理切面:自动记录处理延迟、成功率、积压指标
 */
@Aspect
@Component
public class EventHandlerMetricsAspect {

    private final MeterRegistry registry;

    @Around("@annotation(EventHandler)")
    public Object recordMetrics(ProceedingJoinPoint joinPoint) {
        String eventType = joinPoint.getArgs()[0].getClass().getSimpleName();
        Timer.Sample sample = Timer.start(registry);

        try {
            Object result = joinPoint.proceed();
            registry.counter("event.processed", "event", eventType, "status", "success").increment();
            return result;
        } catch (Throwable e) {
            registry.counter("event.processed", "event", eventType, "status", "failure").increment();
            throw e;
        } finally {
            sample.stop(registry.timer("event.process.duration", "event", eventType));
        }
    }
}

十一、完整集成示例:订单履约事件流

下面是一个贯穿全文的综合运用,展示 EDA + Saga + CQRS + Kafka 的完整集成。

用户下单
    └─> OrderService: 写入 EventStore + Outbox
        └─> Debezium / Poller: 读取 Outbox -> Kafka "order-events"
            ├─> InventoryConsumer: 扣库存 -> 发布 "inventory-reserved"
            ├─> PaymentConsumer:   扣款   -> 发布 "payment-completed"
            ├─> ShippingConsumer:  创建运单 -> 发布 "shipment-created"
            └─> ProjectionService: 更新 CQRS 读模型(ES / Redis / PG)
                └─> Query API: 提供订单列表 / 详情查询

11.1 事件定义汇总

package com.example.eda.integration;

// 所有事件使用密封接口统一管理,便于 switch 模式匹配
public sealed interface OrderDomainEvent permits
    OrderCreatedEvent,
    InventoryReservedEvent,
    InventoryReservationFailedEvent,
    PaymentCompletedEvent,
    PaymentFailedEvent,
    OrderShippedEvent,
    SagaFailedEvent {}

public record OrderCreatedEvent(String orderId, List<OrderLine> lines, Instant occurredOn)
    implements OrderDomainEvent {}

public record InventoryReservedEvent(String orderId, String sku, int qty, Instant occurredOn)
    implements OrderDomainEvent {}

public record PaymentCompletedEvent(String orderId, String transactionId, Instant occurredOn)
    implements OrderDomainEvent {}

public record OrderShippedEvent(String orderId, String trackingNumber, Instant occurredOn)
    implements OrderDomainEvent {}

11.2 基于密封接口的事件路由器

package com.example.eda.integration;

/**
 * 消费者入口:利用 Java 17+ sealed interface + switch 表达式做类型安全路由
 */
@Component
public class OrderEventRouter {

    private final InventoryChoreographyHandler inventoryHandler;
    private final PaymentChoreographyHandler paymentHandler;
    private final OrderListProjection projection;

    public void route(EventEnvelope<OrderDomainEvent> envelope) {
        OrderDomainEvent event = envelope.payload();

        // 先更新读模型投影(尽力而为,失败不影响主流程)
        try {
            switch (event) {
                case OrderCreatedEvent e -> projection.on(e);
                case PaymentCompletedEvent e -> projection.on(e);
                case OrderShippedEvent e -> projection.on(e);
                default -> { /* 投影无需处理的事件 */ }
            }
        } catch (Exception ex) {
            // 投影失败进入死信队列或告警,不阻塞主事件流
            log.error("投影处理失败", ex);
        }

        // 再路由业务处理器
        switch (event) {
            case OrderCreatedEvent e -> inventoryHandler.onOrderCreated(
                new EventEnvelope<>(envelope.eventId(), envelope.correlationId(), envelope.sourceService(), envelope.occurredOn(), e)
            );
            case InventoryReservedEvent e -> paymentHandler.onInventoryReserved(
                new EventEnvelope<>(envelope.eventId(), envelope.correlationId(), envelope.sourceService(), envelope.occurredOn(), e)
            );
            case PaymentCompletedEvent e -> {
                // 触发物流创建
            }
            default -> log.warn("未路由的事件: {}", event.getClass().getSimpleName());
        }
    }
}

十二、FAQ:事件驱动架构高频问题

Q1:事件溯源与一般消息队列有什么区别?

事件溯源强调以事件为唯一真相源,系统状态通过重放事件还原;消息队列仅作为通信管道,下游可独立存储状态,不要求事件长期保留或按序重放。

Q2:Kafka 与 RabbitMQ 在 EDA 中如何选择?

Kafka 胜在高吞吐、持久化、回溯消费与流处理能力,适合事件溯源与日志型事件流。RabbitMQ 胜在低延迟、灵活路由(Exchange-Binding)、AMQP 协议成熟,适合任务队列与 RPC 补齐场景。

Q3:CQRS 一定要配合事件溯源使用吗?

不一定。CQRS 可独立使用,写模型可直接操作关系型数据库,再通过 CDC(如 Debezium)同步到读库。但事件溯源天然生成有序事件流,是 CQRS 投影的最佳事件来源。

Q4:如何处理事件消费者失败导致的无限重试?

建议三层策略:1)业务可重试异常进入指数退避重试队列;2)超过重试阈值转入死信队列(DLQ)人工处理;3)Saga 失败事件触发补偿,自动回滚已完成的步骤。

Q5:微服务拆分后,如何避免因事件泛滥导致系统难以调试?

建立统一事件目录(Event Catalog),每个事件需注册类型名、Schema、版本、生产者与消费方列表。配合 TraceId 全链路追踪,以及事件审计日志(Audit Log)定期审查废弃事件。


总结

事件驱动架构为分布式系统提供了天然的弹性边界与扩展能力。本文从语义层、传输层、存储层、计算层四个维度构建了完整的 EDA 技术栈:

  • 语义层:厘清事件、命令与消息的边界,指导系统设计。
  • 传输层:事件总线实现 Pub-Sub 与队列语义,Kafka 作为高吞吐骨干。
  • 存储层:事件溯源持久化业务事实,快照优化读取性能,CQRS 分离读写模型。
  • 计算层:Saga 协调长事务,Kafka Streams 处理实时流 Join 与窗口聚合。

生产落地时务必关注顺序性、幂等性、Schema 演化与可观测性四大支柱。只有在这些基础设施稳固之后,EDA 才能真正发挥其松耦合、高可用的架构优势,支撑业务的持续演进与规模增长。

继续阅读

探索更多技术文章

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

全部文章 返回首页

「distributed-systems」更多文章

  1. 分布式高可用架构模式:多活、容灾、降级与 K8s 编排高可用
  2. 分布式链路追踪实战:OpenTelemetry、Jaeger 与 W3C Trace Context
  3. 分布式缓存深度策略:Redis Cluster、一致性哈希与多级缓存架构