事件驱动架构(Event-Driven Architecture, EDA)是现代分布式系统的核心通信范式。本文从事件语义拆解出发,逐步深入事件总线、Saga 分布式事务、事件溯源、CQRS、DDD 聚合事件与 Kafka Streams 流编排,提供可直接落地的代码、设计图与生产级实践。
一、EDA 核心语义:事件、消息与命令
在事件驱动系统中,Event、Message、Command 三个概念经常被混用。准确区分它们是设计清晰架构的第一步。
┌──────────────────────────────────────────────────────────────┐
│ 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 Fanout | RabbitMQ 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 才能真正发挥其松耦合、高可用的架构优势,支撑业务的持续演进与规模增长。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。