
Ddd Cqrs Architecture
- 15 installs
- 1 repo stars
- Updated July 29, 2026
- full-statck-skills/ddd-skills
Implements CQRS at L1/L2/L3 levels with Event Sourcing, idempotency, domain-event lifecycle, and integration into five DDD architectures.
About
Guides CQRS adoption across L1/L2/L3 levels including Event Sourcing, idempotency, and outbox-based event projection. A developer uses it when read and write models diverge or high-concurrency writes need eventual-consistency reads.
- Command and query model design
- Outbox, publisher, and projection sync for L2/L3
Ddd Cqrs Architecture by the numbers
- 15 all-time installs (skills.sh)
- Ranked #3,491 of 4,347 Backend & APIs skills by installs in the Skillselion catalog
- Data as of Jul 30, 2026 (Skillselion catalog sync)
npx skills add https://github.com/full-statck-skills/ddd-skills --skill ddd-cqrs-architectureAdd your badge
Show developers this skill is listed on Skillselion. Paste this into your README.
| Installs | 15 |
|---|---|
| repo stars | ★ 1 |
| Last updated | July 29, 2026 |
| Repository | full-statck-skills/ddd-skills ↗ |
What it does
Implements CQRS at L1/L2/L3 levels with Event Sourcing, idempotency, domain-event lifecycle, and integration into five DDD architectures.
Files
DDD CQRS Architecture
CQRS (Command Query Responsibility Segregation) — L1/L2/L3 adoption, Event Sourcing, idempotency, integration with 5 DDD architectures.
When to Use This Skill
Trigger keywords: CQRS, 读写分离, Event Sourcing, 事件溯源, Command Bus, Query Model, 领域事件, domain event, 幂等, idempotency, eventual consistency, projection, materialized view
Workflow
Step 1: 判断是否需要 CQRS → 确认读写模式是否显著分化
Step 2: 选择 CQRS 级别 → L1/L2/L3
Step 3: 设计命令模型 → Command → Handler → 领域事件
Step 4: 设计查询模型 → Query → Handler → DTO
Step 5: 实现事件同步(L2/L3)→ Outbox → 发布器 → 投影
Step 6: 按架构集成 → 根据 Layered/Onion/Hexagonal/Clean/COLAWhen to Use CQRS
适用场景
- 读写模式显著分化: 写操作与读操作用不同数据结构和优化策略
- 高并发写入: 写需 ACID 保证,读可接受最终一致性
- 多视图需求: 同一数据多种展示形式(列表/详情/统计/搜索)
- 审计追踪要求: 完整记录所有状态变更历史
- 团队具备事件驱动能力: 理解最终一致性和事件溯源概念
升级路径
CRUD 够用(读=写) → 单模型,无需 CQRS
↓
L1 模型分离 → CommandService / QueryService 分离,共享 DB
↓
L2 数据库分离 → Command DB + Query DB,事件同步
↓
L3 Event Sourcing → EventStore + Projection不适用场景
| 场景 | 替代方案 |
|---|---|
| 简单 CRUD,读=写 | 单模型,不引入 CQRS |
| 原型/一次性项目 | 跳过 CQRS |
| 团队不熟悉事件驱动 | 先用 ddd-event-storming 建立事件思维 |
| 强一致性要求极高 | 评估分布式事务成本 |
CQRS Core Principles
| 维度 | 命令侧(Write) | 查询侧(Read) |
|---|---|---|
| 职责 | 处理状态变更,执行业务规则 | 返回数据视图,无业务逻辑 |
| 模型 | Command Model(命令对象+聚合根) | Query Model(DTO+物化视图) |
| 存储 | Write DB (3NF, ACID) | Read DB (反范式, 查询优化) |
| 一致性 | 强一致性(聚合内) | 最终一致性(跨聚合/服务) |
| 输出 | 领域事件 | DTO / View Model |
L1/L2/L3 Adoption Strategy
L1 — Model Separation
成本最低:仅代码层分离 Command/Query Service,共享数据库。适用于读写数据结构相同但逻辑分离的场景。 参考示例: examples/order-l1-model-separation.md
L2 — Database Separation
中等成本:分离 Command DB 和 Query DB,通过领域事件同步。适用于读负载高、独立优化策略需求的场景。 参考示例: examples/order-l2-db-separation.md
L3 — Event Sourcing
最高成本:以事件流作为唯一真相源,通过投影重建读模型。适用于审计追踪、时间旅行查询、事件重放的场景。 参考示例: examples/order-l3-event-sourcing.md
Event Lifecycle
领域行为 → 构建 DomainEvent → 持久化 → EventBus 发布
→ 本地处理器 (同步)
→ MQ 外发 (异步, 跨服务)发布策略
- 轮询发布(定时扫表): 延迟 1-5s, 低复杂度, 中小流量
- CDC(Debezium binlog): 延迟 <100ms, 中复杂度, 高流量低延迟
- 事务提交回调: 延迟 <10ms, 低复杂度, Spring 项目
详细参考: references/cqrs-events.md
幂等设计
| 策略 | 开销 | 可靠性 | 场景 |
|---|---|---|---|
| 事件去重表 | 中 | ★★★ | 金融级关键事件 |
| 状态机守卫 | 低 | ★★★ | 状态驱动事件 |
| Redis + TTL | 低 | ★★☆ | 非关键通知 |
| 业务幂等 | 低 | ★★★ | 简单操作 |
详细参考(含代码示例): references/event-governance.md
各架构 CQRS 集成模式
| 架构 | 集成点 | 目录结构 |
|---|---|---|
| Layered | Application 层 Command/Query Service 分离 | app/service/command/ + app/service/query/ |
| Onion | Core 层 Command/Query UseCase 接口 | core/application/command/ + core/application/query/ |
| Hexagonal | Port 层 Command Port + Query Port | domain/port/command/ + domain/port/query/ |
| Clean | UseCase 层 Command Interactor + Query Interactor | usecase/interactor/command/ + usecase/interactor/query/ |
| COLA | App 层 Command + Query 子模块 | app/command/ + app/query/ |
详细参考: examples/multi-architecture-integration.md
Gotchas
- 过早升级到 L2/L3: 先在 L1 验证 CQRS 价值,大多数项目 L1 就够了
- Query 侧直接查写库: Query 绝不能直接读取 Command 侧数据库表
- 忘了事件幂等: 消费者必须实现幂等(at-least-once 投递)
- 事件版本不兼容: 修改事件结构必须向后兼容
- 最终一致性的 UI 处理: L2/L3 下前端需轮询或 WebSocket
- 过大的聚合: 事件溯源聚合应保持小聚合原则
FAQ
| 问题 | 回答 |
|---|---|
| CQRS 一定会引入最终一致性? | L1 不需要,L2/L3 需要 |
| Event Sourcing 和 CQRS 绑定? | 否,可独立使用 |
| 何时需要 Outbox 模式? | 需要可靠事件发布时,避免双写问题 |
| CQRS 和微服务关系? | 正交,可在单体内部署或跨微服务 |
| 查询模型复杂到何种程度? | 可跨聚合/服务组装,反范式化物化视图 |
Rules
| 规则 | 说明 |
|---|---|
| Command 侧需验证权限 | 执行命令前必须鉴权 |
| 事件体不能含敏感数据 | 密码、PII 等不应出现在事件中 |
| 查询侧不能修改状态 | Query Handler/Service 不能有写操作 |
| 同一事务写业务数据 + Outbox | 避免双写问题 |
| 聚合内强一致,聚合间最终一致 | 一个事务最多改一个聚合的状态 |
Skill Boundary
✅ 擅长处理
1. 读写模式显著分化的系统(报表 vs 交易) 2. 需要事件驱动架构的项目 3. CQRS 三级渐进落地(L1 模型分离 / L2 DB分离 / L3 Event Sourcing) 4. 与 5 种 DDD 架构的集成模式 5. 领域事件全生命周期治理(发布/订阅/幂等/补偿)
⚠️ 需要条件
1. 已选好基础架构(Layered/Onion/Hexagonal/Clean/COLA) 2. 团队理解最终一致性和事件驱动概念 3. 业务场景确实需要读写分离(非简单 CRUD)
❌ 不该用(超出范围)
1. 简单 CRUD(读=写) → 不适用 CQRS 2. 原型/一次性项目 → 不应该使用 CQRS 3. 团队不熟悉事件驱动 → 先用 ddd-event-storming 建立事件思维 4. 已使用 CQRS 框架(Axon) → 框架内置了 CQRS 模式
Security & Stability
- Code templates are educational — replace URLs/credentials with env vars.
- Command handlers must validate authorization. Never trust caller is authorized.
- Event payloads must not contain sensitive data — events are stored and replayed.
- 最小权限原则:Command API 和 Query API 应使用不同的访问权限控制
Keywords
CQRS, 读写分离, Event Sourcing, 事件溯源, Command Bus, Query Model, Command Model, 领域事件, Domain Event, 最终一致性, Eventual Consistency, Outbox Pattern, Transactional Outbox, Materialized View, Projection, Event Store, Idempotency, 幂等, L1/L2/L3, Clean Architecture CQRS, Hexagonal CQRS, COLA CQRS
References
Architecture CQRS Patterns
references/architecture/clean-ddd-hexagonal-cqrs.md— CQRS & Domain Events:Commands, Queries, Read Models, Event Dispatcher, Outboxreferences/cqrs-events.md— CQRS 领域事件:事件分类、事件存储、投影策略、版本管理references/cqrs-mindmap.md— CQRS 思维导图:适用场景、实施策略references/ddd4j-cqrs-mindmap.md— CQRS 思维导图:核心理念、架构模式、组件
Domain Events & Event Governance
references/event-governance.md— 事件治理:Outbox DDL、幂等、重试/死信/补偿/对账references/domain-events/domain-events-deep.md— 领域事件深入:事件驱动设计原则references/domain-vs-integration-events.md— 领域事件 vs 集成事件:边界划分references/partme-06-domain-events.md— 领域事件实操:保险承保案例
Examples
examples/order-l1-model-separation.md— L1 模型分离完整示例examples/order-l2-db-separation.md— L2 数据库分离 + 事件同步examples/order-l3-event-sourcing.md— L3 Event Sourcing 完整示例examples/multi-architecture-integration.md— 5 种架构 CQRS 集成对比examples/inventory-cqrs.md— 库存 CQRS + 幂等策略实现
CQRS 测试策略示例
Command/Query 分离后如何编写测试
@SpringBootTest
class OrderCqrsTest {
@Test
void createOrder_Command成功_Query可查() {
// 执行 Command
CreateOrderCommand cmd = new CreateOrderCommand("CUST-1", List.of(
new OrderItemCommand("PROD-1", 2)));
OrderCreatedResult result = commandService.createOrder(cmd);
// 验证 Query
OrderDetailDTO detail = queryService.getOrderDetail(result.getOrderId());
assertThat(detail.getStatus()).isEqualTo("CREATED");
assertThat(detail.getItems()).hasSize(1);
}
@Test
void 事件同步_最终一致验证() {
commandService.payOrder(new PayOrderCommand("ORDER-1"));
// 等待异步事件同步
await().atMost(5, SECONDS).untilAsserted(() -> {
OrderDocument doc = queryService.getOrder("ORDER-1");
assertThat(doc.getStatus()).isEqualTo("PAID");
});
}
}银行账户 Event Sourcing
存款/取款事件溯源,L3 级别,审计追踪
// 事件
public class MoneyDepositedEvent extends DomainEvent {
public final BigDecimal amount;
public final BigDecimal newBalance;
}
public class MoneyWithdrawnEvent extends DomainEvent {
public final BigDecimal amount;
public final BigDecimal newBalance;
}
// 事件溯源聚合
public class BankAccount extends EventSourcedAggregate {
private BigDecimal balance;
public void deposit(BigDecimal amount) {
apply(new MoneyDepositedEvent(this.id, amount, balance.add(amount)));
}
public void withdraw(BigDecimal amount) {
if (balance.compareTo(amount) < 0) throw new InsufficientFundsException();
apply(new MoneyWithdrawnEvent(this.id, amount, balance.subtract(amount)));
}
@Override
protected void when(DomainEvent event) {
if (event instanceof MoneyDepositedEvent e) this.balance = e.newBalance;
else if (event instanceof MoneyWithdrawnEvent e) this.balance = e.newBalance;
}
}
// 投影:账户余额视图
@Component
public class AccountBalanceProjector {
@EventListener
public void on(MoneyDepositedEvent e) {
jdbc.update("UPDATE account_balance SET balance = ? WHERE id = ?",
e.newBalance, e.getAggregateId());
}
}库存 CQRS 完整实现 — 含幂等策略
业务:库存扣减 + 库存查询 | 侧重幂等性和并发控制
业务场景
商品下单 → 扣减库存 → 发布 StockDeductedEvent → 同步到库存视图
幂等设计实现
方案 1: 事件去重表
CREATE TABLE idempotent_event (
event_id VARCHAR(36) PRIMARY KEY,
event_type VARCHAR(100) NOT NULL,
status VARCHAR(20) NOT NULL DEFAULT 'PROCESSING',
handler_name VARCHAR(200) NOT NULL,
created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
processed_at DATETIME
);@Component
public class IdempotentEventHandler {
private final JdbcTemplate jdbc;
public boolean tryProcess(String eventId, String eventType, String handlerName) {
try {
jdbc.update("""
INSERT INTO idempotent_event
(event_id, event_type, status, handler_name)
VALUES (?, ?, 'PROCESSING', ?)
""",
eventId, eventType, handlerName);
return true; // 首次处理
} catch (DuplicateKeyException e) {
// 重复事件,检查状态
String status = jdbc.queryForObject(
"SELECT status FROM idempotent_event WHERE event_id = ?",
String.class, eventId);
return "FAILED".equals(status); // 失败重试
}
}
public void markSuccess(String eventId) {
jdbc.update(
"UPDATE idempotent_event SET status = 'SUCCESS', processed_at = NOW()",
eventId);
}
public void markFailed(String eventId) {
jdbc.update(
"UPDATE idempotent_event SET status = 'FAILED' WHERE event_id = ?",
eventId);
}
}方案 2: 业务条件幂等(库存扣减)
@Component
public class InventoryCommandHandler {
private final JdbcTemplate jdbc;
@Transactional
public DeductResult deductStock(DeductStockCommand command) {
// 幂等扣减: stock >= quantity 时才扣减
int affected = jdbc.update("""
UPDATE inventory
SET stock = stock - ?,
version = version + 1
WHERE product_id = ?
AND stock >= ?
AND version = ?
""",
command.getQuantity(),
command.getProductId(),
command.getQuantity(),
command.getExpectedVersion());
if (affected == 0) {
Inventory current = jdbc.queryForObject(
"SELECT stock, version FROM inventory WHERE product_id = ?",
(rs, row) -> new Inventory(
rs.getInt("stock"),
rs.getInt("version")),
command.getProductId());
if (current == null) {
throw new ProductNotFoundException(command.getProductId());
}
if (current.stock() < command.getQuantity()) {
return DeductResult.insufficient(current.stock());
}
// 并发冲突
return DeductResult.conflict(current.version());
}
return DeductResult.success();
}
}方案 3: 状态机守卫(订单级幂等)
public class InventoryReservation extends AggregateRoot<ReservationId> {
private ReservationStatus status;
private String productId;
private int quantity;
public void confirm() {
if (status == ReservationStatus.CONFIRMED) {
return; // 幂等:已确认则不重复处理
}
if (status != ReservationStatus.PENDING) {
throw new InventoryException("当前状态不可确认: " + status);
}
this.status = ReservationStatus.CONFIRMED;
addDomainEvent(new ReservationConfirmedEvent(this.id));
}
public void cancel() {
if (status == ReservationStatus.CANCELLED) {
return; // 幂等:已取消
}
this.status = ReservationStatus.CANCELLED;
addDomainEvent(new ReservationCancelledEvent(this.id));
}
}完整调用链
// 1. Command 入口
@RestController
@RequestMapping("/api/v1/inventory")
public class InventoryController {
private final InventoryCommandHandler commandHandler;
private final InventoryQueryHandler queryHandler;
@PostMapping("/deduct")
public DeductResult deductStock(@RequestBody DeductStockCommand command) {
return commandHandler.deductStock(command);
}
@GetMapping("/{productId}")
public InventoryView getInventory(@PathVariable String productId) {
return queryHandler.getInventory(productId);
}
}
// 2. 事件发布
@Component
public class StockDeductedEventHandler {
private final IdempotentEventHandler idempotentHandler;
private final InventoryReadRepository readRepo;
@Async
@Transactional
public void on(StockDeductedEvent event) {
String eventId = event.getEventId();
if (!idempotentHandler.tryProcess(eventId, "StockDeductedEvent", "StockDeductedEventHandler")) {
return;
}
try {
// 更新库存视图
readRepo.updateStock(event.getProductId(), event.getRemainingStock());
idempotentHandler.markSuccess(eventId);
} catch (Exception e) {
idempotentHandler.markFailed(eventId);
throw e; // 触发重试
}
}
}
// 3. 查询侧
@Service
public class InventoryQueryHandler {
private final InventoryReadRepository readRepo;
private final RedisTemplate<String, InventoryView> redis;
public InventoryView getInventory(String productId) {
// 缓存提升读性能
return redis.opsForValue()
.get("inventory:" + productId, InventoryView.class)
.orElseGet(() -> {
InventoryView view = readRepo.findByProductId(productId);
redis.opsForValue().set("inventory:" + productId, view, Duration.ofSeconds(30));
return view;
});
}
}幂等策略对比
| 策略 | 代码量 | 性能 | 可靠性 | 适用场景 |
|---|---|---|---|---|
| 去重表 (DB) | 中 | 中 | ★★★ | 跨服务事件,需要持久化去重 |
| 业务条件 (SQL) | 低 | 高 | ★★★ | 库存扣减、余额扣减等数值操作 |
| 状态机守卫 | 低 | 高 | ★★★ | 状态驱动实体(订单、审批) |
| Redis + TTL | 低 | 高 | ★★☆ | 非关键通知、缓存刷新 |
完整测试
@SpringBootTest
class InventoryCqrsTest {
@Test
void deductStock_幂等_同一命令重复执行() {
DeductStockCommand command = new DeductStockCommand("PROD-001", 2, 1);
commandHandler.deductStock(command); // 第一次:成功扣减
DeductResult result = commandHandler.deductStock(command); // 第二次:版本冲突
assertThat(result.isConflict()).isTrue();
}
@Test
void eventHandler_幂等_同一事件重复投递() {
StockDeductedEvent event = new StockDeductedEvent("PROD-001", 98);
eventHandler.on(event); // 第一次:正常处理
eventHandler.on(event); // 第二次:幂等跳过
InventoryView view = queryHandler.getInventory("PROD-001");
assertThat(view.getStock()).isEqualTo(98);
}
}5 种架构 CQRS 集成对比示例
演示同一业务(订单支付)在 5 种 DDD 架构下的 CQRS 集成模式
业务
用户支付订单 → Command: PayOrderCommand → 领域事件: OrderPaidEvent → Query: 订单状态更新1. Layered + CQRS
目录: app/service/command/PayOrderCommandService.java + app/service/query/OrderQueryService.java
// Command Service (Application 层)
@Service
public class PayOrderCommandService {
private final OrderRepository orderRepository;
private final EventPublisher eventPublisher;
@Transactional
public void payOrder(PayOrderCommand command) {
Order order = orderRepository.findById(command.getOrderId())
.orElseThrow(() -> new OrderNotFoundException(command.getOrderId()));
order.pay();
orderRepository.save(order);
eventPublisher.publish(new OrderPaidEvent(order.getId()));
}
}
// Query Service (Application 层)
@Service
public class OrderQueryService {
private final OrderReadRepository readRepository;
public OrderDetailDTO getOrder(String orderId) {
return readRepository.findDetailById(orderId);
}
}2. Onion + CQRS
目录: core/application/command/PayOrderUseCase.java + core/application/query/GetOrderUseCase.java
// 定义 UseCase 接口 (Core 层)
public interface PayOrderUseCase {
void execute(PayOrderCommand command);
}
public interface GetOrderUseCase {
OrderDetailDTO execute(GetOrderQuery query);
}
// 实现 UseCase (Core 层)
public class PayOrderUseCaseImpl implements PayOrderUseCase {
private final OrderRepository orderRepository; // Core 定义的接口
private final EventBus eventBus;
@Override
public void execute(PayOrderCommand command) {
Order order = orderRepository.findById(command.getOrderId());
order.pay();
orderRepository.save(order);
eventBus.publish(new OrderPaidEvent(order.getId()));
}
}
public class GetOrderUseCaseImpl implements GetOrderUseCase {
private final OrderReadRepository readRepo;
@Override
public OrderDetailDTO execute(GetOrderQuery query) {
return readRepo.findDetailById(query.getOrderId());
}
}3. Hexagonal + CQRS
目录: domain/port/command/PayOrderPort.java + domain/port/query/OrderQueryPort.java
// Ports (Domain 层 — 接口定义)
public interface PayOrderPort { // 入站端口
void pay(PayOrderCommand command);
}
public interface OrderQueryPort { // 入站端口
OrderDetailDTO getOrder(String orderId);
}
// UseCase 实现 (Application 层)
@ApplicationService
public class PayOrderService implements PayOrderPort {
private final OrderRepository orderRepo; // 出站端口
private final EventPublisher eventPub; // 出站端口
@Override
public void pay(PayOrderCommand command) {
Order order = orderRepo.findById(new OrderId(command.getOrderId()));
order.pay();
orderRepo.save(order);
eventPub.publish(new OrderPaidEvent(order.getId()));
}
}
// Primary Adapter (REST)
@RestController
public class OrderController {
private final PayOrderPort payOrder; // 注入入站端口
private final OrderQueryPort queryOrder;
@PostMapping("/orders/{id}/pay")
public void payOrder(@PathVariable String id) {
payOrder.pay(new PayOrderCommand(id));
}
@GetMapping("/orders/{id}")
public OrderDetailDTO getOrder(@PathVariable String id) {
return queryOrder.getOrder(id);
}
}4. Clean + CQRS
目录: usecase/interactor/command/PayOrderInteractor.java + usecase/interactor/query/GetOrderInteractor.java
// Input/Output Ports (UseCase 层)
public interface PayOrderInputPort {
void execute(PayOrderInput input);
}
public interface OrderQueryOutputPort {
OrderDetailDTO execute(OrderQueryInput input);
}
// Interactor (UseCase 层)
public class PayOrderInteractor implements PayOrderInputPort {
private final OrderRepository orderRepo; // 输出端口
private final EventPublisher eventPublisher;
@Override
public void execute(PayOrderInput input) {
Order order = orderRepo.findById(input.getOrderId());
order.pay();
orderRepo.save(order);
eventPublisher.publish(new OrderPaidEvent(order.getId()));
}
}
public class GetOrderInteractor implements OrderQueryOutputPort {
private final OrderReadRepository readRepo;
@Override
public OrderDetailDTO execute(OrderQueryInput input) {
return readRepo.findDetailById(input.getOrderId());
}
}5. COLA + CQRS
目录: app/command/OrderCommandExecutor.java + app/query/OrderQueryExecutor.java
// Command 执行器 (App 层)
@Component
public class OrderCommandExecutor {
private final OrderRepository orderRepo;
private final DomainEventBus eventBus;
public void execute(PayOrderCmd cmd) {
Order order = orderRepo.find(cmd.getOrderId());
order.pay();
orderRepo.save(order);
eventBus.fire(new OrderPaidEvent(order.getId()));
}
}
// Query 执行器 (App 层)
@Component
public class OrderQueryExecutor {
private final OrderReadRepo readRepo;
public OrderDetailDTO execute(GetOrderQry qry) {
return readRepo.findDetail(qry.getOrderId());
}
}集成模式对比总结
| 架构 | Command 位置 | Query 位置 | 事件发布者 | 核心抽象 |
|---|---|---|---|---|
| Layered | app/service/command/ | app/service/query/ | 应用服务 | Service 分离 |
| Onion | core/application/command/ | core/application/query/ | UseCase 接口 | UseCase 接口 |
| Hexagonal | domain/port/command/ | domain/port/query/ | 入站端口 | Port |
| Clean | usecase/interactor/command/ | usecase/interactor/query/ | Interactor | Input/Output Port |
| COLA | app/command/ | app/query/ | CommandExecutor | Executor |
通知 CQRS — CQRS 化推送
事件 → 通知投影,L2 级别,Redis 缓冲
// 事件 → 通知创建
@Component
public class OrderNotificationProjector {
@EventListener
public void on(OrderShippedEvent event) {
Notification notif = new Notification(
event.getCustomerId(),
"订单已发货",
"您的订单 " + event.getOrderId() + " 已发出"
);
notificationRepo.save(notif);
}
}
// 通知查询
@Service
public class NotificationQueryService {
public List<NotificationDTO> getUnread(String customerId) {
return notificationRepo.findUnreadByCustomer(customerId);
}
}L1 模型分离完整示例 — Order 聚合
所属级别: L1(模型分离) | 共享同一数据库 | 代码层读写分离
业务场景
用户下单 → 创建订单 → 查询订单列表 / 订单详情
目录结构
com/example/order/
├── command/
│ ├── CreateOrderCommand.java
│ ├── OrderCommandService.java
│ └── OrderCreatedResult.java
├── query/
│ ├── OrderQuery.java
│ ├── OrderQueryService.java
│ ├── OrderDetailDTO.java
│ └── OrderSummaryDTO.java
├── domain/
│ ├── Order.java # 聚合根
│ ├── OrderStatus.java # 值对象
│ ├── OrderId.java # 值对象
│ └── OrderRepository.java # 仓储接口
├── infrastructure/
│ └── OrderRepositoryImpl.java # 仓储实现(共享)
└── controller/
├── OrderCommandController.java
└── OrderQueryController.java写侧代码
Command 对象
public class CreateOrderCommand {
private final String customerId;
private final List<OrderItemCommand> items;
// constructor, getters — immutable
}
public class OrderItemCommand {
private final String productId;
private final int quantity;
}Command Service
@Service
public class OrderCommandService {
private final OrderRepository orderRepository;
@Transactional
public OrderCreatedResult createOrder(CreateOrderCommand command) {
Order order = Order.create(command);
orderRepository.save(order);
return OrderCreatedResult.from(order);
}
}聚合根
public class Order extends AggregateRoot<OrderId> {
private OrderId id;
private OrderStatus status;
private List<OrderItem> items;
private CustomerId customerId;
public static Order create(CreateOrderCommand command) {
Order order = new Order();
order.id = new OrderId(UUID.randomUUID().toString());
order.status = OrderStatus.CREATED;
order.items = command.getItems().stream()
.map(OrderItem::fromCommand)
.toList();
order.customerId = new CustomerId(command.getCustomerId());
order.addDomainEvent(new OrderCreatedEvent(order.id));
return order;
}
}读侧代码
Query 对象
public class OrderQuery {
private final String customerId;
private final OrderStatus status;
private final LocalDate startDate;
private final LocalDate endDate;
}Query Service
@Service
public class OrderQueryService {
private final OrderRepository orderRepository;
public OrderDetailDTO getOrderDetail(OrderId id) {
Order order = orderRepository.findById(id)
.orElseThrow(() -> new OrderNotFoundException(id));
return OrderDetailDTO.from(order);
}
public Page<OrderSummaryDTO> listOrders(OrderQuery query, Pageable pageable) {
return orderRepository.findByCriteria(query, pageable);
}
}DTO
public class OrderDetailDTO {
private String orderId;
private String status;
private String customerName;
private List<OrderItemDTO> items;
private BigDecimal totalAmount;
private LocalDateTime createdAt;
public static OrderDetailDTO from(Order order) {
OrderDetailDTO dto = new OrderDetailDTO();
dto.orderId = order.getId().getValue();
dto.status = order.getStatus().name();
dto.items = order.getItems().stream()
.map(OrderItemDTO::from)
.toList();
dto.createdAt = order.getCreatedAt();
return dto;
}
}控制器
@RestController
@RequestMapping("/api/v1/orders")
public class OrderCommandController {
private final OrderCommandService commandService;
@PostMapping
public ResponseEntity<OrderCreatedResult> createOrder(
@RequestBody CreateOrderCommand command) {
OrderCreatedResult result = commandService.createOrder(command);
return ResponseEntity.status(201).body(result);
}
}
@RestController
@RequestMapping("/api/v1/orders")
public class OrderQueryController {
private final OrderQueryService queryService;
@GetMapping("/{id}")
public ResponseEntity<OrderDetailDTO> getOrder(@PathVariable String id) {
OrderDetailDTO dto = queryService.getOrderDetail(new OrderId(id));
return ResponseEntity.ok(dto);
}
@GetMapping
public ResponseEntity<Page<OrderSummaryDTO>> listOrders(OrderQuery query, Pageable pageable) {
return ResponseEntity.ok(queryService.listOrders(query, pageable));
}
}适用场景
- 读写逻辑已有明显差异,但数据库结构一致
- 最小化架构改动,快速验证 CQRS 价值
- 团队首次引入 CQRS,逐步过渡
L2 数据库分离完整示例 — 事件同步
所属级别: L2(数据库分离) | Command DB + Query DB | 领域事件同步
业务场景
订单支付 → 发布 OrderPaidEvent → 异步同步到 Elasticsearch 查询库
架构图
POST /orders/{id}/pay
↓
OrderCommandService
↓
Order.pay() → OrderPaidEvent → Outbox 表
↓
Outbox Poller / CDC
↓
MQ (Kafka/RabbitMQ)
↓
OrderPaidEventHandler → ES / Read DB
↓
OrderQueryService ← Elasticsearch / Read DB写侧
领域事件
public class OrderPaidEvent extends DomainEvent {
private final String orderId;
private final BigDecimal totalAmount;
private final LocalDateTime paidAt;
public OrderPaidEvent(String orderId, BigDecimal totalAmount) {
super(orderId);
this.orderId = orderId;
this.totalAmount = totalAmount;
this.paidAt = LocalDateTime.now();
}
}聚合行为
public class Order extends AggregateRoot<OrderId> {
private OrderStatus status;
public void pay() {
if (this.status != OrderStatus.PENDING_PAYMENT) {
throw new OrderException("订单状态不可支付: " + this.status);
}
this.status = OrderStatus.PAID;
addDomainEvent(new OrderPaidEvent(this.id.getValue(), this.totalAmount));
}
}Outbox 写入(同一事务)
@Service
public class OrderCommandService {
private final OrderRepository orderRepository;
private final OutboxRepository outboxRepository;
@Transactional
public void payOrder(PayOrderCommand command) {
Order order = orderRepository.findById(new OrderId(command.getOrderId()))
.orElseThrow(() -> new OrderNotFoundException(command.getOrderId()));
order.pay();
orderRepository.save(order);
// 领域事件由仓储自动写入 outbox 表(同一事务)
}
}读侧
事件处理器 → 写查询库
@Component
public class OrderPaidEventHandler {
private final OrderReadRepository readRepo;
@Async
@Transactional
public void handleOrderPaidEvent(OrderPaidEvent event) {
// 幂等检查
if (readRepo.existsById(event.getOrderId())) {
return; // 已处理
}
OrderDocument doc = new OrderDocument();
doc.setOrderId(event.getOrderId());
doc.setStatus("PAID");
doc.setTotalAmount(event.getTotalAmount());
doc.setPaidAt(event.getPaidAt());
doc.setUpdatedAt(LocalDateTime.now());
readRepo.save(doc);
}
}查询服务
@Service
public class OrderQueryService {
private final OrderReadRepository readRepo;
private final ElasticsearchRestTemplate esTemplate;
public OrderDocument getOrder(String orderId) {
return readRepo.findById(orderId)
.orElseThrow(() -> new OrderNotFoundException(orderId));
}
public Page<OrderDocument> search(OrderSearchQuery query, Pageable pageable) {
NativeSearchQuery searchQuery = new NativeSearchQueryBuilder()
.withQuery(QueryBuilders.boolQuery()
.filter(QueryBuilders.termQuery("customerId", query.getCustomerId()))
.filter(QueryBuilders.rangeQuery("paidAt")
.gte(query.getStartDate()).lte(query.getEndDate())))
.withPageable(pageable)
.build();
return esTemplate.search(searchQuery, OrderDocument.class)
.map(SearchHit::getContent);
}
}同步延迟监控
@Component
public class SyncLagMonitor {
private final JdbcTemplate writeDb;
private final JdbcTemplate readDb;
@Scheduled(fixedRate = 60000)
public void checkSyncLag() {
// 对比写库最新事件时间和读库最新记录时间
Instant latestWrite = writeDb.queryForObject(
"SELECT MAX(occurred_at) FROM domain_event_outbox", Instant.class);
Instant latestRead = readDb.queryForObject(
"SELECT MAX(updated_at) FROM order_projection", Instant.class);
Duration lag = Duration.between(latestRead, latestWrite);
if (lag.getSeconds() > 30) {
// 告警:同步延迟超过 30 秒
}
}
}适用场景
- 读负载大,需要独立的查询优化策略
- Command DB (MySQL/PostgreSQL) + Query DB (Elasticsearch/MongoDB)
- 可以接受秒级最终一致性
- 需要复杂的搜索功能(全文检索、聚合统计)
L3 Event Sourcing 完整示例
所属级别: L3(Event Sourcing) | EventStore + Projection | 完整审计追踪
业务场景
订单创建 → 支付 → 发货 → 完成 — 以事件流作为唯一真相源
架构图
Command → Aggregate → EventStream → EventStore (Append-Only)
↓
┌─────┴─────┐
↓ ↓
Projector A Projector B
↓ ↓
OrderView OrderSearch
(MySQL) (ES)Event Store
public interface EventStore {
void appendEvents(
String aggregateId,
List<DomainEvent> events,
int expectedVersion
);
List<DomainEvent> loadEvents(String aggregateId);
List<DomainEvent> loadEventsByType(
String eventType,
LocalDateTime since
);
}
@Repository
public class JdbcEventStore implements EventStore {
private final JdbcTemplate jdbc;
private final ObjectMapper objectMapper;
@Override
public void appendEvents(
String aggregateId,
List<DomainEvent> events,
int expectedVersion) {
for (DomainEvent event : events) {
int inserted = jdbc.update("""
INSERT INTO events (event_id, aggregate_id, aggregate_type,
event_type, event_data, version, occurred_at)
SELECT ?, ?, ?, ?, ?, ?, ?
WHERE (SELECT MAX(version) FROM events
WHERE aggregate_id = ?) = ?
""",
event.getEventId(), aggregateId, "Order",
event.getClass().getSimpleName(),
toJson(event), expectedVersion + 1, event.getOccurredAt(),
aggregateId, expectedVersion
);
if (inserted == 0) {
throw new ConcurrencyConflictException(aggregateId, expectedVersion);
}
}
}
@Override
public List<DomainEvent> loadEvents(String aggregateId) {
return jdbc.query("""
SELECT * FROM events
WHERE aggregate_id = ?
ORDER BY version ASC
""",
new EventRowMapper(), aggregateId
);
}
}Event-Sourced Aggregate
public abstract class EventSourcedAggregate {
private String id;
private int version;
private final List<DomainEvent> pendingEvents = new ArrayList<>();
protected void apply(DomainEvent event) {
pendingEvents.add(event);
when(event); // 应用事件改变状态
}
protected abstract void when(DomainEvent event);
public void loadFromHistory(List<DomainEvent> events) {
events.forEach(e -> {
when(e);
this.version = e.getVersion();
});
}
public List<DomainEvent> getPendingEvents() {
return Collections.unmodifiableList(pendingEvents);
}
public int getVersion() { return version; }
public String getId() { return id; }
}
public class Order extends EventSourcedAggregate {
private OrderStatus status;
private Money totalAmount;
private CustomerId customerId;
public static Order create(CreateOrderCommand cmd) {
Order order = new Order();
order.apply(new OrderCreatedEvent(
UUID.randomUUID().toString(), cmd.getCustomerId(), cmd.getItems()));
return order;
}
public void pay() {
apply(new OrderPaidEvent(this.id));
}
public void ship(TrackingNumber tracking) {
if (status != OrderStatus.PAID) {
throw new OrderException("已支付才能发货");
}
apply(new OrderShippedEvent(this.id, tracking));
}
@Override
protected void when(DomainEvent event) {
if (event instanceof OrderCreatedEvent e) {
this.id = e.getAggregateId();
this.status = OrderStatus.CREATED;
this.customerId = new CustomerId(e.getCustomerId());
} else if (event instanceof OrderPaidEvent) {
this.status = OrderStatus.PAID;
} else if (event instanceof OrderShippedEvent) {
this.status = OrderStatus.SHIPPED;
}
}
}Command Handler
@Service
public class OrderCommandHandler {
private final EventStore eventStore;
@Transactional
public void handlePayOrder(PayOrderCommand command) {
List<DomainEvent> events = eventStore.loadEvents(command.getOrderId());
Order order = new Order();
order.loadFromHistory(events);
order.pay();
eventStore.appendEvents(
command.getOrderId(),
order.getPendingEvents(),
order.getVersion()
);
}
}Projection
@Component
public class OrderProjection {
private final JdbcTemplate jdbc;
@EventListener
public void onOrderCreated(OrderCreatedEvent event) {
jdbc.update("""
INSERT INTO order_projection
(order_id, status, customer_id, total_amount, created_at)
VALUES (?, 'CREATED', ?, ?, ?)
""",
event.getAggregateId(), event.getCustomerId(),
event.getTotalAmount(), event.getOccurredAt());
}
@EventListener
public void onOrderPaid(OrderPaidEvent event) {
jdbc.update(
"UPDATE order_projection SET status = 'PAID', paid_at = ? WHERE order_id = ?",
event.getOccurredAt(), event.getAggregateId());
}
@EventListener
public void onOrderShipped(OrderShippedEvent event) {
jdbc.update(
"UPDATE order_projection SET status = 'SHIPPED', tracking_no = ? WHERE order_id = ?",
event.getTrackingNumber().getValue(), event.getAggregateId());
}
}Snapshot 策略
@Component
public class OrderSnapshotter {
private final EventStore eventStore;
private final JdbcTemplate jdbc;
// 每 100 个事件创建快照
private static final int SNAPSHOT_INTERVAL = 100;
@Scheduled(fixedRate = 3600000)
public void createSnapshot(String aggregateId) {
List<DomainEvent> events = eventStore.loadEvents(aggregateId);
if (events.size() % SNAPSHOT_INTERVAL == 0) {
Order order = new Order();
order.loadFromHistory(events);
jdbc.update("""
INSERT INTO aggregate_snapshots
(aggregate_id, aggregate_type, snapshot_data, version, created_at)
VALUES (?, 'Order', ?, ?, ?)
""",
aggregateId, serialize(order), order.getVersion(), Instant.now());
}
}
}Events DDL
CREATE TABLE events (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
event_id VARCHAR(36) NOT NULL UNIQUE,
aggregate_id VARCHAR(36) NOT NULL,
aggregate_type VARCHAR(100) NOT NULL,
event_type VARCHAR(200) NOT NULL,
event_data JSON NOT NULL,
version INT NOT NULL,
occurred_at DATETIME NOT NULL,
created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
UNIQUE KEY uk_aggregate_version (aggregate_id, version),
INDEX idx_events_type_time (event_type, occurred_at)
);
CREATE TABLE aggregate_snapshots (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
aggregate_id VARCHAR(36) NOT NULL,
aggregate_type VARCHAR(100) NOT NULL,
snapshot_data JSON NOT NULL,
version INT NOT NULL,
created_at DATETIME NOT NULL,
INDEX idx_snapshot_aggregate (aggregate_id, version)
);适用场景
- 需要完整审计追踪(金融、合规)
- 时间旅行查询("某时间点的状态是什么")
- 事件重放(从零重建系统状态)
- 事件驱动分析(事件流 → 数据分析)
支付 CQRS — 事件驱动对账
支付命令 + 支付事件 → 对账视图,L2 级别
// 支付命令
public class ProcessPaymentCommand {
private final String orderId;
private final BigDecimal amount;
}
// 支付事件 → 对账投影
@Component
public class PaymentEventHandler {
@EventListener
public void on(PaymentCompletedEvent event) {
// 更新对账视图
reconciliationRepo.recordPayment(
event.getOrderId(), event.getAmount(), event.getPaidAt());
}
}
// 对账查询
@Service
public class ReconciliationQueryService {
public List<ReconciliationDTO> getDailyReport(LocalDate date) {
return reconciliationRepo.findByDate(date);
}
}用户 CQRS — 命令与查询分离
用户注册(Command) + 用户查询(Query),L1 级别
// Command
public class RegisterUserCommand {
private final String username;
private final String email;
private final String password;
}
@Service
public class UserCommandService {
public UserRegisteredResult register(RegisterUserCommand cmd) {
User user = User.register(cmd);
userRepository.save(user);
return UserRegisteredResult.from(user);
}
}
// Query
@Service
public class UserQueryService {
public UserProfileDTO getProfile(String userId) {
return userReadRepo.findProfileById(userId);
}
}CQRS & Domain Events
Sources:
- CQRS — Martin Fowler
- Event Sourcing — Martin Fowler
- CQRS Pattern — Microsoft Azure
- Transactional Outbox — microservices.io
- Domain Events – Salvation — Udi Dahan
- Strengthening Your Domain: Domain Events — Jimmy Bogard
- Domain Events: Design and Implementation — Microsoft
CQRS Overview
Command Query Responsibility Segregation separates read and write operations into different models.
flowchart TB
API["API Layer"]
API --> Commands
API --> Queries
subgraph WriteSide["Write Side"]
Commands["Commands"]
CmdHandler["Command Handler\n(Use Case)"]
DomainModel["Domain Model\n(Aggregates)"]
WriteDB[("Write Database")]
Commands --> CmdHandler
CmdHandler --> DomainModel
DomainModel --> WriteDB
end
subgraph ReadSide["Read Side"]
Queries["Queries"]
QryHandler["Query Handler\n(Read Model)"]
ReadDB[("Read Database\n(Optimized)")]
Queries --> QryHandler
QryHandler --> ReadDB
end
WriteDB -->|Domain Events| EventHandler["Event Handler"]
EventHandler -->|Updates| ReadDB
style WriteSide fill:#3b82f6,stroke:#2563eb,color:white
style ReadSide fill:#10b981,stroke:#059669,color:white
style EventHandler fill:#f59e0b,stroke:#d97706,color:white---
Commands vs Queries
Commands (Write Side)
Commands represent intent to change state. They mutate data.
// application/commands/place_order_command.ts
export interface PlaceOrderCommand {
type: 'PlaceOrder';
customerId: string;
items: Array<{
productId: string;
quantity: number;
}>;
}
export interface ConfirmOrderCommand {
type: 'ConfirmOrder';
orderId: string;
}
export interface CancelOrderCommand {
type: 'CancelOrder';
orderId: string;
reason: string;
}
export class PlaceOrderHandler {
async handle(command: PlaceOrderCommand): Promise<OrderId> {
const order = Order.create(CustomerId.from(command.customerId));
for (const item of command.items) {
const product = await this.productRepo.findById(item.productId);
order.addItem(product.id, item.quantity, product.price);
}
await this.orderRepo.save(order);
await this.eventPublisher.publishAll(order.domainEvents);
return order.id;
}
}Queries (Read Side)
Queries retrieve data without side effects. They never mutate state.
// application/queries/get_order_query.ts
export interface GetOrderQuery {
orderId: string;
}
export interface GetOrdersByCustomerQuery {
customerId: string;
status?: OrderStatus;
page?: number;
pageSize?: number;
}
export interface OrderDTO {
id: string;
customerId: string;
customerName: string;
status: string;
items: Array<{
productId: string;
productName: string;
quantity: number;
unitPrice: number;
subtotal: number;
}>;
total: number;
createdAt: string;
confirmedAt?: string;
}
export class GetOrderHandler {
constructor(private readonly readDb: IOrderReadModel) {}
async handle(query: GetOrderQuery): Promise<OrderDTO | null> {
return this.readDb.findById(query.orderId);
}
}
export class GetOrdersByCustomerHandler {
constructor(private readonly readDb: IOrderReadModel) {}
async handle(query: GetOrdersByCustomerQuery): Promise<PaginatedResult<OrderDTO>> {
return this.readDb.findByCustomer(
query.customerId,
query.status,
query.page ?? 1,
query.pageSize ?? 20
);
}
}---
Read Model (Projection)
Optimized database structure for queries. Can denormalize data for performance.
interface IOrderReadModel:
findById(orderId: string) -> OrderDTO | null
findByCustomer(customerId, status?, page?, pageSize?) -> PaginatedResult<OrderDTO>
search(criteria: OrderSearchCriteria) -> List<OrderDTO>
class PostgresOrderReadModel implements IOrderReadModel:
db: Database
findById(orderId: string) -> OrderDTO | null:
row = db.ordersRead
.where(id: orderId)
.join("customer")
.withRelated("items.product")
.first()
return row ? this.mapToDTO(row) : nullSeparate write and read databases (optional): write is normalized for transactions, read is denormalized for queries.
---
Domain Events
Notifications that something happened in the domain. Used for:
- Updating read models
- Cross-aggregate communication
- Integration with other bounded contexts
Event Structure
// domain/shared/domain_event.ts
export abstract class DomainEvent {
readonly eventId: string;
readonly occurredAt: Date;
readonly aggregateId: string;
abstract readonly eventType: string;
constructor(aggregateId: string) {
this.eventId = crypto.randomUUID();
this.occurredAt = new Date();
this.aggregateId = aggregateId;
}
abstract toPayload(): Record<string, unknown>;
}
// domain/order/events.ts
export class OrderCreated extends DomainEvent {
readonly eventType = 'order.created';
constructor(
readonly orderId: OrderId,
readonly customerId: CustomerId,
) {
super(orderId.value);
}
toPayload() {
return {
orderId: this.orderId.value,
customerId: this.customerId.value,
};
}
}
export class OrderConfirmed extends DomainEvent {
readonly eventType = 'order.confirmed';
constructor(
readonly orderId: OrderId,
readonly total: Money,
readonly items: ReadonlyArray<{ productId: string; quantity: number }>,
) {
super(orderId.value);
}
toPayload() {
return {
orderId: this.orderId.value,
total: { amount: this.total.amount, currency: this.total.currency },
items: this.items,
};
}
}
export class OrderShipped extends DomainEvent {
readonly eventType = 'order.shipped';
constructor(
readonly orderId: OrderId,
readonly trackingNumber: string,
readonly carrier: string,
) {
super(orderId.value);
}
toPayload() {
return {
orderId: this.orderId.value,
trackingNumber: this.trackingNumber,
carrier: this.carrier,
};
}
}Event Handlers
class OrderCreatedHandler:
db: Database
handle(event: OrderCreated):
db.ordersRead.insert({
id: event.orderId.value,
customerId: event.customerId.value,
status: "draft",
createdAt: event.occurredAt
})
class OrderConfirmedHandler:
db: Database
handle(event: OrderConfirmed):
db.ordersRead
.where(id: event.orderId.value)
.update({
status: "confirmed",
total: event.total.amount,
confirmedAt: event.occurredAt
})
export class SendShippingNotificationHandler {
constructor(
private readonly orderRepo: IOrderRepository,
private readonly notifier: INotificationService,
) {}
async handle(event: OrderShipped): Promise<void> {
const order = await this.orderRepo.findById(OrderId.from(event.orderId.value));
if (!order) return;
await this.notifier.sendEmail(order.customerEmail, {
template: 'order-shipped',
data: {
orderId: event.orderId.value,
trackingNumber: event.trackingNumber,
carrier: event.carrier,
},
});
}
}---
Domain Events vs Integration Events
Domain Events
- Stay within bounded context
- Fine-grained, low-level
- Trigger internal processes
- Named in domain language
class OrderItemQuantityIncreased extends DomainEvent {
constructor(
readonly orderId: OrderId,
readonly productId: ProductId,
readonly oldQuantity: number,
readonly newQuantity: number,
) { super(orderId.value); }
}Integration Events
- Cross bounded context boundaries
- Coarser-grained
- Published to message broker
- Versioned schema
interface OrderConfirmedIntegrationEvent {
eventType: 'sales.order.confirmed';
eventId: string;
version: '1.0';
occurredAt: string;
payload: {
orderId: string;
customerId: string;
total: { amount: number; currency: string };
items: Array<{
productId: string;
quantity: number;
unitPrice: number;
}>;
shippingAddress: {
street: string;
city: string;
postalCode: string;
country: string;
};
};
}Publishing Integration Events
// application/event_handlers/publish_integration_events.ts
export class PublishOrderConfirmedIntegrationEvent {
constructor(
private readonly messageBroker: IMessageBroker,
private readonly orderRepo: IOrderRepository,
) {}
async handle(domainEvent: OrderConfirmed): Promise<void> {
const order = await this.orderRepo.findById(domainEvent.orderId);
if (!order) return;
const integrationEvent: OrderConfirmedIntegrationEvent = {
eventType: 'sales.order.confirmed',
eventId: crypto.randomUUID(),
version: '1.0',
occurredAt: new Date().toISOString(),
payload: {
orderId: order.id.value,
customerId: order.customerId.value,
total: {
amount: order.total.amount,
currency: order.total.currency,
},
items: order.items.map(item => ({
productId: item.productId.value,
quantity: item.quantity.value,
unitPrice: item.unitPrice.amount,
})),
shippingAddress: order.shippingAddress
? {
street: order.shippingAddress.street,
city: order.shippingAddress.city,
postalCode: order.shippingAddress.postalCode,
country: order.shippingAddress.country,
}
: null,
},
};
await this.messageBroker.publish('order-events', integrationEvent);
}
}---
Event Dispatcher Pattern
// infrastructure/events/event_dispatcher.ts
export interface IEventHandler<T extends DomainEvent> {
handle(event: T): Promise<void>;
}
export class EventDispatcher {
private handlers: Map<string, IEventHandler<any>[]> = new Map();
register<T extends DomainEvent>(
eventType: string,
handler: IEventHandler<T>,
): void {
const existing = this.handlers.get(eventType) ?? [];
existing.push(handler);
this.handlers.set(eventType, existing);
}
async dispatch(event: DomainEvent): Promise<void> {
const handlers = this.handlers.get(event.eventType) ?? [];
await Promise.all(handlers.map(h => h.handle(event)));
}
async dispatchAll(events: DomainEvent[]): Promise<void> {
for (const event of events) {
await this.dispatch(event);
}
}
}
const dispatcher = new EventDispatcher();
dispatcher.register('order.created', new OrderCreatedHandler(readDb));
dispatcher.register('order.confirmed', new OrderConfirmedHandler(readDb));
dispatcher.register('order.confirmed', new PublishOrderConfirmedIntegrationEvent(broker, orderRepo));
dispatcher.register('order.shipped', new SendShippingNotificationHandler(orderRepo, notifier));---
Outbox Pattern
Ensures events are published reliably (exactly-once semantics).
interface OutboxMessage:
id: string
eventType: string
payload: string
createdAt: DateTime
processedAt: DateTime | null
class OutboxRepository:
db: Database
save(event: DomainEvent, tx: Transaction):
tx.outbox.insert({
id: event.eventId,
eventType: event.eventType,
payload: serialize(event.toPayload()),
createdAt: event.occurredAt
})
getUnprocessed(limit: int = 100) -> List<OutboxMessage>:
return db.outbox
.where(processedAt: null)
.orderBy("createdAt")
.limit(limit)
.lockForUpdate()
markProcessed(id: string):
db.outbox.where(id: id).update({processedAt: now()})
class PlaceOrderHandler:
orderRepo: IOrderRepository
outbox: OutboxRepository
db: Database
handle(command: PlaceOrderCommand) -> OrderId:
order = Order.create(CustomerId.from(command.customerId))
for item in command.items:
product = productRepo.findById(item.productId)
order.addItem(product.id, item.quantity, product.price)
db.transaction(tx => {
orderRepo.save(order, tx)
for event in order.domainEvents:
outbox.save(event, tx)
})
return order.id---
Event Sourcing (Brief Overview)
Store state changes as a sequence of events rather than current state.
Event Store:
Stream: order-123
Events:
1. OrderCreated{customerId: "c1", items: [...]}
2. OrderItemAdded{productId: "p1", quantity: 2}
3. OrderConfirmed{total: $40}
4. OrderShipped{trackingNumber: "TN123"}
Current State = fold over all events (left fold)Use when:
- Complete audit trail is required
- Temporal queries needed ("what was state on June 1?")
- Event replay for debugging/testing
Skip when:
- Simple CRUD operations
- No audit requirements
- Team unfamiliar with pattern
CQRS 实现路径:L1 → L2 → L3 渐进式演进
CQRS 不是一次性全部引入,而是根据业务复杂度逐步升级。每个 Level 都有明确的判断条件和可逆回退方案。
CQRS 三级实现速查
| Level | 读写模型 | 数据库 | 同步方式 | 复杂度 | 适用场景 |
|---|---|---|---|---|---|
| L1 | 代码层分离(Command/Query DTO) | 同一 DB | 同步 | 低 | 读写模型差异小 |
| L2 | 物理分离(Read Model 独立表) | 同一 DB | 异步(事件) | 中 | 查询需要 JOIN 优化 |
| L3 | 完全分离(独立读库) | 不同 DB | 异步(事件+CDC) | 高 | 读写性能要求差距大 / Event Sourcing |
---
L1:代码层分离(Command/Query DTO 分离)
判断何时升级到 L1
当前状态:所有操作使用同一 Service + 同一 DTO
→ 读返回不必要的字段、写校验混在读模型
触发条件(满足 2 项即可升级):
- 读接口字段数 > 写接口字段数 × 2
- 同一实体有 3+ 种不同的查询视图
- 写操作响应中包含大量冗余关联数据第 1 步:拆分 Command 和 Query Service
// 改造前(混在一起)
@Service
public class OrderService {
public OrderDTO createOrder(CreateOrderRequest req) { ... }
public OrderDTO getOrder(String id) { ... }
public List<OrderDTO> listOrders(OrderQuery query) { ... }
}
// 改造后(L1:代码分离)
@Service
public class OrderCommandService {
public OrderId createOrder(CreateOrderCommand cmd) {
Order order = Order.create(cmd.getCustomerId(), cmd.getItems());
orderRepository.save(order);
return order.getId();
}
}
@Service
public class OrderQueryService {
public OrderDetailDTO getOrder(String id) {
return orderReadModel.findById(id);
}
public PageResult<OrderSummaryDTO> listOrders(OrderQuery query) {
return orderReadModel.findByCriteria(query);
}
}第 2 步:定义独立的 Query DTO
// 不同查询场景使用不同 DTO
public class OrderDetailDTO { // 详情页:完整信息
private String orderId;
private String status;
private String totalAmount;
private List<OrderItemDTO> items;
private String customerName;
private String shippingAddress;
}
public class OrderSummaryDTO { // 列表页:核心字段
private String orderId;
private String status;
private String totalAmount;
private LocalDateTime createdAt;
}第 3 步:验证 L1 完成
# 检查:Command Service 是否只返回 ID(不返回 DTO)
grep -r "return.*DTO" **/command/ # 应无结果
# 检查:Query DTO 是否不作为写操作的响应
grep -r "QueryService" **/controller/ | grep "POST\|PUT\|DELETE" # 应无结果L1 常见错误
| 错误 | 后果 | 修复 |
|---|---|---|
| Command Handler 返回完整 DTO | 写操作耦合读模型 | 返回 ID + 状态码 |
| QueryService 中包含写操作 | 读写未真正分离 | 拆分到 CommandService |
| 所有查询共用一个 DTO | DTO 膨胀 | 按查询场景定义独立 DTO |
---
L2:读模型物理分离(同一 DB,独立表/视图)
判断何时升级到 L2
L1 已实施,且满足以下条件之一:
- 查询需要 3+ 表 JOIN,响应超过 200ms
- 不同查询视图 > 5 种,L1 DTO 维护困难
- 需要数据聚合(SUM/COUNT/GROUP BY)但写模型结构不支持第 1 步:设计 Read Model 表
-- 写模型(规范化)
CREATE TABLE orders (
id VARCHAR(36) PRIMARY KEY,
customer_id VARCHAR(36) NOT NULL,
status VARCHAR(20) NOT NULL,
created_at TIMESTAMP DEFAULT NOW(),
version INT DEFAULT 0
);
-- 读模型(反规范化,专为查询优化)
CREATE TABLE order_read_model (
order_id VARCHAR(36) PRIMARY KEY,
customer_id VARCHAR(36),
customer_name VARCHAR(100), -- JOIN 结果平铺
status VARCHAR(20),
total_amount DECIMAL(12,2), -- 预计算
item_count INT, -- 预计算
created_at TIMESTAMP,
updated_at TIMESTAMP,
INDEX idx_customer (customer_id),
INDEX idx_status_created (status, created_at)
);第 2 步:实现 Event Handler 同步 Read Model
@Component
public class OrderReadModelProjector {
private final OrderReadModelMapper readModelMapper;
@EventListener
@Transactional
public void on(OrderCreatedEvent event) {
OrderReadModel model = OrderReadModel.builder()
.orderId(event.getOrderId())
.customerId(event.getCustomerId())
.status("DRAFT")
.createdAt(event.getOccurredAt())
.build();
readModelMapper.insert(model);
}
@EventListener
@Transactional
public void on(OrderPaidEvent event) {
readModelMapper.updateStatus(event.getOrderId(), "PAID");
}
}第 3 步:Query Service 直接查 Read Model
@Mapper
public interface OrderReadModelMapper {
OrderDetailDTO findById(@Param("orderId") String orderId);
List<OrderSummaryDTO> findByCustomer(
@Param("customerId") String customerId,
@Param("status") String status,
@Param("offset") int offset,
@Param("limit") int limit
);
}第 4 步:验证 L2 完成
# 检查:写操作不直接查询 Read Model
grep -r "ReadModel" **/command/ # 应无结果
# 检查:Event Handler 覆盖率
# 每个关键业务操作都应有对应的 Event Handler 更新 Read ModelL2 常见错误
| 错误 | 后果 | 修复 |
|---|---|---|
| 跳过 Outbox 直接发事件 | 事件丢失导致读写不一致 | 事务内写 Outbox |
| Event Handler 耗时过长 | 写事务阻塞 | Handler 异步化,不能放在写事务中 |
| Read Model 只增不改 | 旧数据未被清理 | 加 updated_at 字段和清理策略 |
---
L3:独立读库(不同 DB 或 ES/Cache)
判断何时升级到 L3
L2 已实施,且满足以下条件之一:
- 读 QPS > 写 QPS × 10,且读写竞争同一 DB 连接池
- 需要全文搜索、实时聚合(ES 场景)
- 读模型数据量 > 1000 万行,MySQL 查询变慢
- 需要 Event Sourcing(审计、回放、时间旅行)第 1 步:配置独立数据源
@Configuration
public class ReadDataSourceConfig {
@Bean
@ConfigurationProperties("spring.datasource.read")
public DataSource readDataSource() {
return DataSourceBuilder.create().build();
}
@Bean
public SqlSessionFactory readSqlSessionFactory() throws Exception {
SqlSessionFactoryBean bean = new SqlSessionFactoryBean();
bean.setDataSource(readDataSource());
return bean.getObject();
}
}第 2 步:同步策略选择
| 同步方式 | 延迟 | 一致性 | 复杂度 | 适用 |
|---|---|---|---|---|
| CDC (Debezium) | < 100ms | 最终一致 | 中 | MySQL → ES / PG |
| Outbox + 消息队列 | 1-5s | 最终一致 | 中 | 跨服务同步 |
| 双写(事务内) | 0ms | 强一致 | 高 | 要求实时一致 |
// CDC 方式(Debezium 监听 MySQL binlog → Kafka → Consumer 写入 ES)
@Component
public class OrderSyncedToElasticsearch {
@KafkaListener(topics = "cdc.orders.order_read_model")
public void syncToES(String orderJson) {
OrderDoc doc = objectMapper.readValue(orderJson, OrderDoc.class);
elasticsearchTemplate.index(doc);
}
}第 3 步:验证 L3 完成
# 检查:写操作依赖只包含写库
grep -r "readDataSource\|readDb" **/command/ # 应无结果
# 检查:读库同步延迟监控
# 监控指标:read_model.updated_at - orders.updated_at 的 p99 < 5s
# 检查:对账机制已配置
# 定时任务检测:order_read_model 中缺失的记录L3 常见错误
| 错误 | 后果 | 修复 |
|---|---|---|
| 跳过延迟监控 | 读库滞后数小时无人知 | 对账任务 + data_lag 指标告警 |
| CDC 误删数据 | 读库历史数据丢失 | CDC 用 INSERT+DELETE 日志,不用 TRUNCATE |
| 写库故障影响读 | 全站不可用 | 读服务有读库兜底(缓存/静态数据) |
---
逆向回退方案
如果 CQRS 引入后复杂度超出预期,可按以下路径回退:
L3 → L2:停用独立读库,Query Service 回退到主库的 Read Model 表
L2 → L1:删除 Read Model 表,Query Service 直接查主库
L1 → 无:合并 Command/Query Service,恢复单一 OrderService
回退检查清单:
- [ ] Read Model 数据是否需要备份
- [ ] 消费者是否依赖独立读库的连接
- [ ] 回退后性能是否可接受(压力测试验证)CQRS 思维导图 — 知识体系参考
核心理念
CQRS = 读写分离
├── 命令侧(写模型)
│ ├── 处理状态变更
│ ├── 执行业务规则
│ └── 产生领域事件
└── 查询侧(读模型)
├── 返回数据视图
├── 无业务逻辑
└── 高度优化分离程度级别
| 级别 | 描述 | 复杂度 |
|---|---|---|
| L1:代码分离 | Command + Query 代码路径分离,共用 DB | 低 |
| L2:数据库分离 | 写库/读库分离,事件同步 | 中 |
| L3:服务分离 | 独立微服务,独立部署 | 高 |
| L4:物理分离 | 独立数据中心/区域 | 极高 |
命令侧组件
- 命令对象:DTO 模式,意图明确,不可变
- 命令处理器:处理单个命令,调用领域逻辑,事务管理
- 聚合根:业务规则封装,状态变更,事件产生
- 事件发布器:领域事件 → 集成事件 → 事件存储
查询侧组件
- 查询处理器:无业务逻辑,纯数据检索,DTO 组装
- 数据模型:物化视图、非规范化、查询优化
- 数据同步策略:事件订阅、CDC、轮询、手动同步
- 存储选项:关系型DB / NoSQL(MongoDB, ES, Redis) / 搜索引擎
事件处理系统
事件类型:领域事件 / 集成事件 / 系统事件
消息队列:RabbitMQ / Kafka / AWS SNS/SQS
事件处理器:同步/异步/竞争消费者/幂等处理
投影机制:实时投影/批量投影/重放/版本处理一致性模型
| 类型 | 特点 |
|---|---|
| 强一致性 | 立即更新,复杂实现,不常用 |
| 最终一致性 | 延迟更新,事件传播,补偿机制 |
| 冲突解决 | 乐观并发、版本控制、冲突检测 |
适用场景判断
| 强烈推荐 | 可以考虑 | 不推荐 |
|---|---|---|
| 高并发写入 | 协作域 | 简单 CRUD |
| 复杂业务规则 | 复杂查询需求 | 强一致性要求高 |
| 不同读写负载 | 事件溯源需求 | 团队经验不足 |
| 需要审计追踪 | 性能优化 | 小规模应用 |
| 多视图需求 | — | — |
渐进式实施策略
单体内先分离代码路径 → 再分离数据库 → 再独立服务
Don't skip levels.源代码
CQRS(命令查询职责分离)思维导图
📌 核心理念
├── 基本概念
│ ├── 读写分离原则
│ ├── 命令侧(写模型)
│ │ ├── 处理状态变更
│ │ ├── 执行业务规则
│ │ └── 产生领域事件
│ └── 查询侧(读模型)
│ ├── 返回数据视图
│ ├── 无业务逻辑
│ └── 高度优化
├── 核心分离原则
│ ├── 职责分离
│ ├── 数据模型分离
│ ├── 数据存储分离(可选)
│ └── 代码路径分离
└── 分离程度级别
├── 级别1:代码分离
├── 级别2:数据库分离
├── 级别3:服务分离
└── 级别4:物理分离🏗️ 架构模式
├── 基础CQRS
│ ├── 单一数据库
│ ├── 共享业务逻辑
│ └── 简单分离
├── 高级CQRS + ES
│ ├── 事件溯源
│ ├── 事件存储
│ ├── 投影/物化视图
│ └── 最终一致性
└── 混合模式
├── CQRS + DDD
├── CQRS + 微服务
└── CQRS + 六边形架构🔧 命令侧(写模型)
├── 命令处理流程
│ ├── 接收命令
│ ├── 验证命令
│ ├── 执行业务逻辑
│ ├── 更新聚合状态
│ ├── 发布领域事件
│ └── 返回结果
├── 组件
│ ├── 命令对象
│ │ ├── DTO模式
│ │ ├── 意图明确
│ │ └── 不可变性
│ ├── 命令处理器
│ │ ├── 处理单个命令
│ │ ├── 调用领域逻辑
│ │ └── 事务管理
│ ├── 聚合根
│ │ ├── 业务规则封装
│ │ ├── 状态变更
│ │ └── 事件产生
│ └── 事件发布器
│ ├── 领域事件
│ ├── 集成事件
│ └── 事件存储
└── 存储选项
├── 关系型数据库
├── 事件存储
└── 文档数据库🔍 查询侧(读模型)
├── 查询处理流程
│ ├── 接收查询
│ ├── 检索数据
│ ├── 格式化响应
│ └── 返回DTO
├── 组件
│ ├── 查询对象
│ │ ├── 只读参数
│ │ └── 无副作用
│ ├── 查询处理器
│ │ ├── 无业务逻辑
│ │ ├── 数据检索
│ │ └── DTO组装
│ └── 数据模型
│ ├── 物化视图
│ ├── 非规范化
│ └── 查询优化
├── 数据同步策略
│ ├── 事件订阅
│ ├── 变更数据捕获
│ ├── 轮询
│ └── 手动同步
└── 存储选项
├── 关系型数据库
├── NoSQL数据库
│ ├── MongoDB
│ ├── Elasticsearch
│ └── Redis
├── 内存数据库
└── 搜索引擎🔄 事件处理系统
├── 事件类型
│ ├── 领域事件
│ ├── 集成事件
│ └── 系统事件
├── 事件总线/消息队列
│ ├── RabbitMQ
│ ├── Kafka
│ ├── Azure Service Bus
│ └── AWS SNS/SQS
├── 事件处理器
│ ├── 同步处理器
│ ├── 异步处理器
│ ├── 竞争消费者
│ └── 幂等性处理
└── 投影/物化视图构建
├── 实时投影
├── 批量投影
├── 重放机制
└── 版本处理⚖️ 一致性模型
├── 强一致性(不常用)
│ ├── 立即更新
│ └── 复杂实现
├── 最终一致性
│ ├── 延迟更新
│ ├── 事件传播
│ └── 补偿机制
├── 数据新鲜度
│ ├── 近实时
│ ├── 延迟容忍
│ └── 陈旧数据策略
└── 冲突解决
├── 乐观并发
├── 版本控制
├── 冲突检测
└── 解决策略🛡️ 挑战与解决方案
├── 复杂性增加
│ ├── 解方案:渐进采用
│ └── 解方案:充分评估
├── 数据同步延迟
│ ├── 解方案:明确预期
│ └── 解方案:监控告警
├── 事件顺序保证
│ ├── 解方案:顺序ID
│ ├── 解方案:分区键
│ └── 解方案:因果关系
├── 错误处理
│ ├── 重试机制
│ ├── 死信队列
│ └── 手动干预
└── 测试复杂性
├── 单元测试
├── 集成测试
└── 端到端测试🎯 适用场景
├── 强烈推荐场景
│ ├── 高并发写入
│ ├── 复杂业务规则
│ ├── 不同读写负载
│ ├── 需要审计追踪
│ └── 多视图需求
├── 可以考虑场景
│ ├── 协作域
│ ├── 复杂查询需求
│ ├── 事件溯源需求
│ └── 性能优化
└── 不推荐场景
├── 简单CRUD
├── 强一致性要求高
├── 团队经验不足
└── 小规模应用📈 实施策略
├── 渐进式实施
│ ├── 从单体开始
│ ├── 先分离代码
│ ├── 再分离数据库
│ └── 最后分离服务
├── 团队准备
│ ├── 技能培训
│ ├── 概念理解
│ └── 模式实践
└── 监控与运维
├── 延迟监控
├── 事件流监控
├── 数据一致性检查
└── 健康状况检查🔧 技术栈示例
├── .NET生态
│ ├── MediatR
│ ├── NServiceBus
│ ├── EventStore
│ └── Marten
├── Java生态
│ ├── Axon Framework
│ ├── Spring Cloud Stream
│ ├── Kafka Streams
│ └── JPA/Hibernate
├── JavaScript/TypeScript
│ ├── NestJS CQRS模块
│ ├── EventEmitter
│ └── TypeORM
└── 云原生方案
├── AWS EventBridge + DynamoDB
├── Azure Functions + Cosmos DB
└── Google Pub/Sub + Firestore✅ 最佳实践
1. 明确分离边界 - 清晰的命令/查询划分 2. 从简单开始 - 避免过早优化 3. 事件设计 - 语义化、版本化 4. 幂等处理 - 确保可靠性 5. 监控先行 - 建立可观测性 6. 文档完善 - 记录数据流和约定 7. 回滚策略 - 准备好故障恢复
---
思维导图使用建议:
- 中心主题:CQRS
- 主要分支:理念、架构、命令侧、查询侧、事件、一致性
- 用颜色区分复杂程度(红:高难度,绿:基础)
- 添加实际案例参考
- 标注适用性和警告提示
这个思维导图可以帮助团队理解CQRS的全貌,并根据实际需求选择适合的实现程度。
领域事件深度参考
领域事件识别方法
捕捉业务专家口中的关键词:
- "如果发生……,则……"
- "当做完……的时候,请通知……"
- "发生……时,则……"
---
事件设计 5 步法
Step 1: 从业务流程中提取事件
业务描述 → 过去式事件
"用户提交订单后,系统扣减库存,然后等待支付"
→ 订单已提交(OrderSubmitted)
→ 库存已扣减(InventoryDeducted)
"支付成功后,订单状态变更为已支付,发送通知"
→ 支付已完成(PaymentCompleted)
→ 订单已支付(OrderPaid)
→ 通知已发送(NotificationSent)Step 2: 确定事件上下文
public class OrderPaid extends DomainEvent {
private final OrderId orderId;
private final Money paidAmount;
private final PaymentMethod paymentMethod;
private final LocalDateTime paidAt;
// 必须包含足够消费方使用的数据
// 但不包含敏感信息(信用卡号等)
}Step 3: 确定发布策略
事件发布位置决策:
├── 同一事务内(聚合内强一致)
│ └── order.pay() 后 addDomainEvent(OrderPaid)
├── 同一进程内微服务内
│ └── Outbox 表 + 定时轮询发布
└── 跨微服务
└── Outbox → MQ(Kafka/RabbitMQ)Step 4: 设计消费者幂等
@EventListener
@Transactional
public void onOrderPaid(OrderPaidEvent event) {
// Step 1: 幂等检查
if (eventRecordDao.exists(event.getEventId())) {
log.info("Duplicate event ignored: {}", event.getEventId());
return;
}
// Step 2: 业务处理
orderReadModel.updateStatus(event.getOrderId(), "PAID");
// Step 3: 记录消费
eventRecordDao.insert(new EventRecord(event.getEventId()));
}Step 5: 事件版本管理
// V1 原始字段
class OrderPaid extends DomainEvent {
private OrderId orderId;
private Money amount;
}
// V2 新增字段(向后兼容)
class OrderPaid extends DomainEvent {
private OrderId orderId;
private Money amount;
private PaymentMethod paymentMethod; // V2 新增,旧消费者忽略
}---
事件持久化实现(Outbox 模式)
-- 与业务表在同一数据库,同一事务写入
CREATE TABLE domain_event_outbox (
id VARCHAR(36) PRIMARY KEY,
aggregate_id VARCHAR(36) NOT NULL,
event_type VARCHAR(200) NOT NULL,
event_data JSONB NOT NULL,
occurred_at TIMESTAMP NOT NULL,
published BOOLEAN NOT NULL DEFAULT FALSE,
retry_count INTEGER NOT NULL DEFAULT 0,
INDEX idx_unpublished (published, created_at)
);@Transactional
public void payOrder(PayOrderCommand cmd) {
Order order = orderRepo.findById(cmd.getOrderId()).orElseThrow();
order.pay(); // 执行业务
orderRepo.save(order); // 持久化聚合
for (DomainEvent event : order.getDomainEvents()) {
outboxRepo.save(OutboxRecord.from(event)); // 同一事务写 Outbox
}
}---
发布策略对比
| 策略 | 延迟 | 复杂度 | 实现 |
|---|---|---|---|
| 轮询发布 | 1-5s | 低 | @Scheduled(fixedDelay=2000) 扫 Outbox 表 |
| CDC (Debezium) | < 100ms | 中 | 监听 binlog → Kafka → Consumer |
| 事务提交回调 | < 10ms | 低 | TransactionSynchronization.afterCommit() |
---
承保业务流程案例
投保微服务 → 生成缴费通知单事件 → 收款微服务订阅并缴费
收款微服务 → 缴费已完成事件 → 投保微服务转保单
投保微服务 → 保单已生成事件 → 保单微服务保存
保单微服务 → 事件扇出 → 佣金/收付费/再保/财务等微服务---
为什么用最终一致性
- 一次事务最多只能更改一个聚合的状态
- 涉及多个聚合状态变更时,用领域事件实现最终一致
- 切断领域模型之间的强依赖关系
- 实现 1 个发布方 → N 个订阅方
---
常见错误
| 错误 | 后果 | 修复 |
|---|---|---|
| 事件不含足够上下文 | 消费者需要反查数据库 | 事件体包含消费方需要的全部数据 |
| 事件和业务不在同一事务 | 业务成功但事件丢失 | Outbox 模式,事务内写入 |
| 消费者不幂等 | 重复消费导致数据错误 | 事件去重表 + UNIQUE 约束 |
| 事件版本不兼容 | 升级后消费者解析失败 | 只增字段,不删不改 |
领域事件 vs 集成事件 + 事件分发器模式
领域事件(Domain Events)
- 限于界上下文内部
- 细粒度、底层
- 触发内部流程
- 用领域语言命名
public class OrderItemQuantityIncreased extends DomainEvent {
private final OrderId orderId;
private final ProductId productId;
private final int oldQuantity;
private final int newQuantity;
public OrderItemQuantityIncreased(OrderId orderId, ProductId productId,
int oldQuantity, int newQuantity) {
super(orderId.getValue());
this.orderId = orderId;
this.productId = productId;
this.oldQuantity = oldQuantity;
this.newQuantity = newQuantity;
}
}集成事件(Integration Events)
- 跨限界上下文
- 粗粒度
- 发布到消息中间件
- 版本化 schema
public class OrderConfirmedIntegrationEvent {
private final String eventType = "sales.order.confirmed";
private final String eventId;
private final String version = "1.0";
private final Instant occurredAt;
private final OrderConfirmedPayload payload;
public static OrderConfirmedIntegrationEvent from(Order order) {
return new OrderConfirmedIntegrationEvent(
UUID.randomUUID().toString(), Instant.now(),
new OrderConfirmedPayload(
order.getId().getValue(),
order.getCustomerId().getValue(),
order.getTotal(),
order.getItems(),
order.getShippingAddress()
)
);
}
}发布集成事件
@Component
public class PublishOrderConfirmedIntegrationEvent {
private final MessageBroker broker;
private final OrderRepository orderRepo;
@EventListener
public void handle(OrderConfirmed domainEvent) {
Order order = orderRepo.findById(domainEvent.getOrderId()).orElse(null);
if (order == null) return;
var integrationEvent = OrderConfirmedIntegrationEvent.from(order);
broker.publish("order-events", integrationEvent);
}
}事件分发器(Event Dispatcher)模式
public class EventDispatcher {
private final Map<String, List<EventHandler<?>>> handlers = new HashMap<>();
public <T extends DomainEvent> void register(String eventType, EventHandler<T> handler) {
handlers.computeIfAbsent(eventType, k -> new ArrayList<>()).add(handler);
}
public void dispatch(DomainEvent event) {
handlers.getOrDefault(event.getEventType(), List.of())
.forEach(h -> h.handle(event));
}
public void dispatchAll(List<DomainEvent> events) {
events.forEach(this::dispatch);
}
}
// 注册
var dispatcher = new EventDispatcher();
dispatcher.register("order.created", new OrderCreatedHandler());
dispatcher.register("order.confirmed", new OrderConfirmedHandler());
dispatcher.register("order.confirmed", new PublishOrderConfirmedIntegrationEvent(broker, repo));
dispatcher.register("order.shipped", new SendShippingNotificationHandler(repo, notifier));事件流全貌
领域行为(同一事务)
→ 写业务数据
→ 在聚合上 addDomainEvent()
→ 事务提交后 EventDispatcher.dispatchAll(events)
├── 内部 Handler(同一进程,同步)
└── 集成事件发布器 → MQ(异步)领域事件 vs 集成事件对照
| 维度 | 领域事件 | 集成事件 |
|---|---|---|
| 范围 | 限界上下文内 | 跨限界上下文 |
| 粒度 | 细粒度 | 粗粒度 |
| 传输 | 进程内 EventBus/Spring Events | MQ(Kafka/RabbitMQ) |
| Schema | 内部约定 | 版本化 + 向后兼容 |
| 命名 | OrderItemAdded | sales.order.confirmed |
事件驱动工程化治理
整合自原 ddd-eventing-governance 技能,覆盖 Outbox、幂等、重试死信、补偿对账、事件契约与可观测性。1. 事件处理闭环
领域行为(同一事务)
→ 写业务数据
→ 写 Outbox 事件表
→ 发布器投递(至少一次)
→ 消费者落库去重(幂等)
→ 执行业务处理
→ 失败:重试/死信/补偿/对账2. Outbox 方案
2.1 Outbox 表 DDL
-- PostgreSQL
CREATE TABLE domain_event_outbox (
id VARCHAR(36) PRIMARY KEY,
aggregate_id VARCHAR(36) NOT NULL,
aggregate_type VARCHAR(100) NOT NULL,
event_type VARCHAR(200) NOT NULL,
event_data JSONB NOT NULL,
schema_version INTEGER NOT NULL DEFAULT 1,
occurred_at TIMESTAMP NOT NULL,
published BOOLEAN NOT NULL DEFAULT FALSE,
published_at TIMESTAMP,
retry_count INTEGER NOT NULL DEFAULT 0,
last_error TEXT,
correlation_id VARCHAR(36),
trace_id VARCHAR(36),
created_at TIMESTAMP NOT NULL DEFAULT NOW(),
INDEX idx_outbox_published (published, created_at),
INDEX idx_outbox_aggregate (aggregate_id),
INDEX idx_outbox_type (event_type, occurred_at)
);-- MySQL
CREATE TABLE domain_event_outbox (
id VARCHAR(36) NOT NULL,
aggregate_id VARCHAR(36) NOT NULL,
aggregate_type VARCHAR(100) NOT NULL,
event_type VARCHAR(200) NOT NULL,
event_data JSON NOT NULL,
schema_version INT NOT NULL DEFAULT 1,
occurred_at DATETIME NOT NULL,
published TINYINT NOT NULL DEFAULT 0,
published_at DATETIME,
retry_count INT NOT NULL DEFAULT 0,
last_error TEXT,
correlation_id VARCHAR(36),
trace_id VARCHAR(36),
created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
PRIMARY KEY (id),
INDEX idx_outbox_published (published, created_at),
INDEX idx_outbox_aggregate (aggregate_id)
);2.2 发布策略对比
| 策略 | 机制 | 延迟 | 复杂度 | 适用场景 |
|---|---|---|---|---|
| 轮询发布 | 定时任务扫 Outbox 表 | 1-5s | 低 | 中小流量,简单可靠 |
| CDC (Debezium) | 监听 binlog/WAL 变更 | < 100ms | 中 | 高流量,低延迟要求 |
| 事务提交后回调 | TransactionSynchronization | < 10ms | 低 | Spring 项目首选 |
2.3 轮询发布器实现
@Component
public class OutboxPublisher {
private final DomainEventOutboxMapper outboxMapper;
private final ApplicationEventPublisher eventPublisher;
private final int BATCH_SIZE = 100;
private final int MAX_RETRY = 5;
@Scheduled(fixedDelay = 2000)
public void publishUnpublishedEvents() {
List<EventOutboxRecord> records = outboxMapper
.findUnpublished(BATCH_SIZE);
for (EventOutboxRecord record : records) {
try {
DomainEvent event = deserialize(record);
eventPublisher.publishEvent(event);
outboxMapper.markPublished(record.getId());
} catch (Exception e) {
handlePublishFailure(record, e);
}
}
}
private void handlePublishFailure(EventOutboxRecord record, Exception e) {
if (record.getRetryCount() >= MAX_RETRY) {
outboxMapper.markDead(record.getId(), e.getMessage());
alertDeadLetter(record);
} else {
outboxMapper.incrementRetry(record.getId(), e.getMessage());
}
}
}2.4 事务中写入 Outbox(同一事务保证)
@Service
public class OrderApplicationService {
private final OrderRepository orderRepository;
private final DomainEventOutboxMapper outboxMapper;
@Transactional
public void payOrder(PayOrderCommand cmd) {
// 1. 执行业务操作
Order order = orderRepository.findById(cmd.getOrderId()).orElseThrow();
order.pay();
// 2. 提取领域事件
List<DomainEvent> events = order.getDomainEvents();
// 3. 持久化聚合(业务数据)
orderRepository.save(order);
// 4. 同一事务写入 Outbox
for (DomainEvent event : events) {
EventOutboxRecord record = EventOutboxRecord.from(event);
outboxMapper.insert(record);
}
}
}3. 幂等策略
3.1 消费幂等表 DDL
-- PostgreSQL
CREATE TABLE event_consumption_record (
event_id VARCHAR(36) NOT NULL,
consumer_group VARCHAR(100) NOT NULL,
consumed_at TIMESTAMP NOT NULL DEFAULT NOW(),
status VARCHAR(20) NOT NULL DEFAULT 'PROCESSED',
PRIMARY KEY (event_id, consumer_group)
);-- MySQL
CREATE TABLE event_consumption_record (
event_id VARCHAR(36) NOT NULL,
consumer_group VARCHAR(100) NOT NULL,
consumed_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
status VARCHAR(20) NOT NULL DEFAULT 'PROCESSED',
PRIMARY KEY (event_id, consumer_group)
);3.2 幂等消费模板
@Component
public class IdempotentEventHandler {
private final EventConsumptionRecordMapper recordMapper;
@EventListener
@Transactional
public void handleOrderPaid(OrderPaidEvent event) {
String eventId = event.getEventId();
// 幂等检查:数据库唯一约束防重
if (recordMapper.exists(eventId, "order-consumer")) {
log.info("Duplicate event ignored: {}", eventId);
return;
}
try {
// 执行业务处理
doHandleOrderPaid(event);
// 记录消费
recordMapper.insert(eventId, "order-consumer", "PROCESSED");
} catch (Exception e) {
// 记录失败(不阻挡重试)
recordMapper.insert(eventId, "order-consumer", "FAILED_RETRY");
throw e;
}
}
}3.3 策略对比
| 策略 | 机制 | 可靠性 | 性能 | 适用 |
|---|---|---|---|---|
| 事件表 + 唯一约束 | DB UNIQUE(event_id, consumer) | ★★★ | ★★☆ | 金融、订单等关键业务 |
| 状态机Guard | if alreadyPaid() return | ★★★ | ★★★ | 状态驱动的业务 |
| Redis SETNX + TTL | SETNX event:xxx EX 3600 | ★★☆ | ★★★ | 非关键通知 |
| 业务幂等 | UPDATE stock=stock-1 WHERE stock>=1 | ★★★ | ★★★ | 简单原子操作 |
4. 重试与死信策略
4.1 重试策略配置
| 参数 | 默认值 | 说明 |
|---|---|---|
| 最大重试次数 | 3 | 超过则进入死信 |
| 初始退避 | 1s | 首次重试延迟 |
| 退避倍数 | 2 | 指数退避:1s → 2s → 4s |
| 最大退避 | 60s | 退避上限 |
| 最大重试时长 | 5min | 超时则标记失败 |
| 死信告警 | 是 | 进入死信立即告警 |
4.2 死信策略
死信处理流程:
消费失败 → 重试耗尽 → 写入死信表/死信队列
→ 钉钉/飞书/邮件告警
→ 人工排查(对账工具查询 eventId 全链路)
→ 修复后手动重放 或 标记为跳过
死信表 DDL:
CREATE TABLE event_dead_letter (
id VARCHAR(36) PRIMARY KEY,
event_id VARCHAR(36) NOT NULL,
event_type VARCHAR(200),
consumer_group VARCHAR(100),
error_message TEXT,
error_stack TEXT,
retry_count INT,
dead_at TIMESTAMP DEFAULT NOW(),
status VARCHAR(20) DEFAULT 'PENDING', -- PENDING/RESOLVED/SKIPPED
resolved_at TIMESTAMP,
resolved_by VARCHAR(100),
INDEX idx_dlq_status (status, dead_at)
);4.3 对账机制
@Component
public class EventReconciliationJob {
private final OutboxMapper outboxMapper;
private final ConsumptionRecordMapper consumptionMapper;
private final DeadLetterMapper deadLetterMapper;
@Scheduled(cron = "0 0 3 * * ?") // 每天凌晨3点
public void reconcile() {
// 1. 扫描超过24小时未发布的事件
List<EventOutboxRecord> stuckEvents = outboxMapper
.findStuckEvents(Duration.ofHours(24));
// 2. 扫描已发布但无消费记录的事件
List<EventOutboxRecord> unconsumedEvents = outboxMapper
.findPublishedButNotConsumed();
// 3. 生成对账报告
ReconciliationReport report = ReconciliationReport.builder()
.stuckPublishEvents(stuckEvents)
.unconsumedEvents(unconsumedEvents)
.deadLetterCount(deadLetterMapper.countPending())
.build();
// 4. 差异自动修复(安全操作)
for (EventOutboxRecord event : stuckEvents) {
if (event.getRetryCount() < 10) {
outboxMapper.resetForRetry(event.getId());
}
}
// 5. 发送对账报告
alertService.sendReconciliationReport(report);
}
}5. 事件契约与版本策略
5.1 事件契约字段规范
每个领域事件必须包含以下字段:
{
"eventId": "uuid",
"eventType": "OrderPaid",
"aggregateId": "ORD-2024-001",
"aggregateType": "Order",
"occurredAt": "2024-03-15T10:30:00Z",
"schemaVersion": 2,
"correlationId": "corr-xxx",
"traceId": "trace-xxx",
"payload": {
"orderId": "ORD-2024-001",
"customerId": "CUS-001",
"totalAmount": { "amount": "99.00", "currency": "CNY" },
"items": [...]
}
}5.2 版本兼容策略
| 策略 | 做法 | 适用 |
|---|---|---|
| 只增字段 | 新增字段设默认值,旧消费者忽略新字段 | 兼容升级 |
| 升级主版本 | 破坏性变更发布到新 topic,旧 topic 保留 | API 签名变化 |
| 双写过渡 | 同时发布 v1 和 v2,待消费者升级后下线 v1 | 平滑迁移 |
5.3 版本降级处理
// 消费者兼容多版本
@EventHandler
public void on(OrderPaidEventV2 event) {
// V2 fields
String paymentMethod = event.getPaymentMethod(); // V2 新增
// Legacy V1 fields still work
Money amount = event.getTotalAmount();
}6. 可观测性
6.1 追踪传播
// 生产者端:传递 traceId + correlationId
public class EventOutboxRecord {
public static EventOutboxRecord from(DomainEvent event) {
return EventOutboxRecord.builder()
.traceId(MDC.get("traceId")) // 当前请求 traceId
.correlationId(UUID.randomUUID().toString())
.build();
}
}
// 消费者端:恢复上下文
@EventListener
public void handle(DomainEvent event) {
MDC.put("traceId", event.getTraceId());
MDC.put("correlationId", event.getCorrelationId());
try {
doHandle(event);
} finally {
MDC.clear();
}
}6.2 关键指标
| 指标 | 说明 | 告警阈值 |
|---|---|---|
events.published.rate | 事件发布速率 | — |
events.publish.failures | 发布失败数 | > 10/min |
events.consumed.rate | 消费速率 | — |
events.consume.failures | 消费失败数 | > 5/min |
events.dead_letter.count | 死信积压数 | > 0 |
events.lag.seconds | 端到端延迟 p95 | > 60s |
outbox.stuck.count | 超过5分钟未发布事件 | > 0 |
6.3 事件全链路查询
-- 按 eventId 查询全链路
SELECT
e.id,
e.event_type,
e.occurred_at AS published_time,
c.consumed_at,
CASE
WHEN e.published = FALSE THEN 'STUCK'
WHEN c.event_id IS NULL THEN 'PUBLISHED_NOT_CONSUMED'
WHEN d.id IS NOT NULL THEN 'DEAD_LETTER(' || d.status || ')'
ELSE 'CONSUMED'
END AS status
FROM domain_event_outbox e
LEFT JOIN event_consumption_record c
ON e.id = c.event_id
LEFT JOIN event_dead_letter d
ON e.id = d.event_id
WHERE e.aggregate_id = :aggregateId
ORDER BY e.occurred_at;7. 验收清单
- [ ] Outbox 表已创建,与业务表在同一数据库
- [ ] 业务操作与 Outbox 写入在同一事务
- [ ] 发布策略已选定(轮询/CDC/回调)并实现
- [ ] 消费幂等表已创建,有唯一约束
- [ ] 所有消费者实现幂等检查
- [ ] 重试次数、退避策略已配置
- [ ] 死信表/队列已创建
- [ ] 死信告警已配置(钉钉/飞书/邮件)
- [ ] 对账任务已配置(定时扫描未消费事件)
- [ ] 事件契约包含 eventId/aggregateId/schemaVersion/correlationId/traceId
- [ ] 版本兼容策略已明确(只增字段/双写/新 topic)
- [ ] traceId 在生产者和消费者间正确传播
- [ ] 关键指标已接入监控(发布延迟/失败率/死信积压)
- [ ] 事件全链路查询可用(按 aggregateId 或 eventId)
在事件风暴(Event Storming)时,我们发现除了命令和操作等业务行为以外,还有一种非常重要的事件,
这种事件发生后通常会导致进一步的业务操作,在 DDD 中这种事件被称为领域事件。
这只是最简单的定义,并不能让我们真正理解它。那到底什么是领域事件?领域事件的技术实现机制是怎样的?这一讲,我们就重点解决这两个大的问题。
领域事件是领域模型中非常重要的一部分,用来表示领域中发生的事件。一个领域事件将导致进一步的业务操作,在实现业务解耦的同时,还有助于形成完整的业务闭环。
举例来说的话,领域事件可以是业务流程的一个步骤,比如投保业务缴费完成后,触发投保单转保单的动作;也可能是定时批处理过程中发生的事件,比如批处理生成季缴保费通知单,触发发送缴费邮件通知操作;或者一个事件发生后触发的后续动作,比如密码连续输错三次,触发锁定账户的动作。
那如何识别领域事件呢?
很简单,和刚才讲的定义是强关联的。在做用户旅程或者场景分析时,我们要捕捉业务、需求人员或领域专家口中的关键词:“如果发生……,则……”“当做完……的时候,请通知……”“发生……时,则……”等。在这些场景中,如果发生某种事件后,会触发进一步的操作,那么这个事件很可能就是领域事件。
那领域事件为什么要用最终一致性,而不是传统 SOA 的直接调用的方式呢?
我们一起回顾一下
讲到的聚合的一个设计原则:
在边界之外使用最终一致性。
一次事务最多只能更改一个聚合的状态。如果一次业务操作涉及多个聚合状态的更改,应采用领域事件的最终一致性。
领域事件驱动设计可以切断领域模型之间的强依赖关系,事件发布完成后,发布方不必关心后续订阅方事件处理是否成功,这样可以实现领域模型的解耦,维护领域模型的独立性和数据的一致性。在领域模型映射到微服务系统架构时,领域事件可以解耦微服务,微服务之间的数据不必要求强一致性,而是基于事件的最终一致性。
回到具体的业务场景,我们发现有的领域事件发生在微服务内的聚合之间,有的则发生在微服务之间,还有两者皆有的场景,一般来说跨微服务的领域事件处理居多。在微服务设计时不同领域事件的处理方式会不一样。
1. 微服务内的领域事件
当领域事件发生在微服务内的聚合之间,领域事件发生后完成事件实体构建和事件数据持久化,发布方聚合将事件发布到事件总线,订阅方接收事件数据完成后续业务操作。
微服务内大部分事件的集成,都发生在同一个进程内,进程自身可以很好地控制事务,因此不一定需要引入消息中间件。但一个事件如果同时更新多个聚合,按照 DDD“一次事务只更新一个聚合”的原则,你就要考虑是否引入事件总线。但微服务内的事件总线,可能会增加开发的复杂度,因此你需要结合应用复杂度和收益进行综合考虑。
微服务内应用服务,可以通过跨聚合的服务编排和组合,以服务调用的方式完成跨聚合的访问,这种方式通常应用于实时性和数据一致性要求高的场景。这个过程会用到分布式事务,以保证发布方和订阅方的数据同时更新成功。
2. 微服务之间的领域事件
跨微服务的领域事件会在不同的限界上下文或领域模型之间实现业务协作,其主要目的是实现微服务解耦,减轻微服务之间实时服务访问的压力。
领域事件发生在微服务之间的场景比较多,事件处理的机制也更加复杂。跨微服务的事件可以推动业务流程或者数据在不同的子域或微服务间直接流转。
跨微服务的事件机制要总体考虑事件构建、发布和订阅、事件数据持久化、消息中间件,甚至事件数据持久化时还可能需要考虑引入分布式事务机制等。
微服务之间的访问也可以采用应用服务直接调用的方式,实现数据和服务的实时访问,弊端就是跨微服务的数据同时变更需要引入分布式事务,以确保数据的一致性。分布式事务机制会影响系统性能,增加微服务之间的耦合,所以我们还是要尽量避免使用分布式事务。
领域事件相关案例
我来给你介绍一个保险承保业务过程中有关领域事件的案例。
一个保单的生成,经历了很多子域、业务状态变更和跨微服务业务数据的传递。这个过程会产生很多的领域事件,这些领域事件促成了保险业务数据、对象在不同的微服务和子域之间的流转和角色转换。
在下面这张图中,我列出了几个关键流程,用来说明如何用领域事件驱动设计来驱动承保业务流程。
事件起点:客户购买保险 - 业务人员完成保单录入 - 生成投保单 - 启动缴费动作。
1. 投保微服务生成缴费通知单,发布第一个事件:缴费通知单已生成,将缴费通知单数据发布到消息中间件。收款微服务订阅缴费通知单事件,完成缴费操作。缴费通知单已生成,领域事件结束。
2. 收款微服务缴费完成后,发布第二个领域事件:缴费已完成,将缴费数据发布到消息中间件。原来的订阅方收款微服务这时则变成了发布方。原来的事件发布方投保微服务转换为订阅方。投保微服务在收到缴费信息并确认缴费完成后,完成投保单转成保单的操作。缴费已完成,领域事件结束。
3. 投保微服务在投保单转保单完成后,发布第三个领域事件:保单已生成,将保单数据发布到消息中间件。保单微服务接收到保单数据后,完成保单数据保存操作。保单已生成,领域事件结束。
4. 保单微服务完成保单数据保存后,后面还会发生一系列的领域事件,以并发的方式将保单数据通过消息中间件发送到佣金、收付费和再保等微服务,一直到财务,完后保单后续所有业务流程。这里就不详细说了。
总之,通过领域事件驱动的异步化机制,可以推动业务流程和数据在各个不同微服务之间的流转,实现微服务的解耦,减轻微服务之间服务调用的压力,提升用户体验。
领域事件总体架构
领域事件的执行需要一系列的组件和技术来支撑。我们来看一下这个领域事件总体技术架构图,
领域事件处理包括:事件构建和发布、事件数据持久化、事件总线、消息中间件、事件接收和处理等。
下面我们逐一讲一下。
1. 事件构建和发布
事件基本属性至少包括:事件唯一标识、发生时间、事件类型和事件源,其中事件唯一标识应该是全局唯一的,以便事件能够无歧义地在多个限界上下文中传递。事件基本属性主要记录事件自身以及事件发生背景的数据。
另外事件中还有一项更重要,那就是业务属性,用于记录事件发生那一刻的业务数据,这些数据会随事件传输到订阅方,以开展下一步的业务操作。
事件基本属性和业务属性一起构成事件实体,事件实体依赖聚合根。领域事件发生后,事件中的业务数据不再修改,因此业务数据可以以序列化值对象的形式保存,这种存储格式在消息中间件中也比较容易解析和获取。
为了保证事件结构的统一,我们还会创建事件基类 DomainEvent(参考下图),子类可以扩充属性和方法。由于事件没有太多的业务行为,实现方法一般比较简单。
事件发布之前需要先构建事件实体并持久化。事件发布的方式有很多种,你可以通过应用服务或者领域服务发布到事件总线或者消息中间件,也可以从事件表中利用定时程序或数据库日志捕获技术获取增量事件数据,发布到消息中间件。
2. 事件数据持久化
事件数据持久化可用于系统之间的数据对账,或者实现发布方和订阅方事件数据的审计。当遇到消息中间件、订阅方系统宕机或者网络中断,在问题解决后仍可继续后续业务流转,保证数据的一致性。
事件数据持久化有两种方案,在实施过程中你可以根据自己的业务场景进行选择。
持久化到本地业务数据库的事件表中,利用本地事务保证业务和事件数据的一致性。
持久化到共享的事件数据库中。这里需要注意的是:业务数据库和事件数据库不在一个数据库中,它们的数据持久化操作会跨数据库,因此需要分布式事务机制来保证业务和事件数据的强一致性,结果就是会对系统性能造成一定的影响。
3. 事件总线 (EventBus)
事件总线是实现微服务内聚合之间领域事件的重要组件,它提供事件分发和接收等服务。事件总线是进程内模型,它会在微服务内聚合之间遍历订阅者列表,采取同步或异步的模式传递数据。事件分发流程大致如下:
如果是微服务内的订阅者(其它聚合),则直接分发到指定订阅者;
如果是微服务外的订阅者,将事件数据保存到事件库(表)并异步发送到消息中间件;
如果同时存在微服务内和外订阅者,则先分发到内部订阅者,将事件消息保存到事件库(表),再异步发送到消息中间件。
4. 消息中间件
跨微服务的领域事件大多会用到消息中间件,实现跨微服务的事件发布和订阅。消息中间件的产品非常成熟,市场上可选的技术也非常多,比如 Kafka,RabbitMQ 等。
5. 事件接收和处理
微服务订阅方在应用层采用监听机制,接收消息队列中的事件数据,完成事件数据的持久化后,就可以开始进一步的业务处理。领域事件处理可在领域服务中实现。
领域事件运行机制相关案例
这里我用承保业务流程的缴费通知单事件,来给你解释一下领域事件的运行机制。这个领域事件发生在投保和收款微服务之间。发生的领域事件是:缴费通知单已生成。下一步的业务操作是:缴费。
事件起点:出单员生成投保单,核保通过后,发起生成缴费通知单的操作。
1. 投保微服务应用服务,调用聚合中的领域服务 createPaymentNotice 和 createPaymentNoticeEvent,分别创建缴费通知单、缴费通知单事件。其中缴费通知单事件类 PaymentNoticeEvent 继承基类 DomainEvent。
2. 利用仓储服务持久化缴费通知单相关的业务和事件数据。为了避免分布式事务,这些业务和事件数据都持久化到本地投保微服务数据库中。
3. 通过数据库日志捕获技术或者定时程序,从数据库事件表中获取事件增量数据,发布到消息中间件。这里说明:事件发布也可以通过应用服务或者领域服务完成发布。
4. 收款微服务在应用层从消息中间件订阅缴费通知单事件消息主题,监听并获取事件数据后,应用服务调用领域层的领域服务将事件数据持久化到本地数据库中。
5. 收款微服务调用领域层的领域服务 PayPremium,完成缴费。
6. 事件结束。
提示:缴费完成后,后续流程的微服务还会产生很多新的领域事件,比如缴费已完成、保单已保存等等。这些后续的事件处理基本上跟 1~6 的处理机制类似。
今天我们主要讲了领域事件以及领域事件的处理机制。领域事件驱动是很成熟的技术,在很多分布式架构中得到了大量的使用。领域事件是 DDD 的一个重要概念,在设计时我们要重点关注领域事件,用领域事件来驱动业务的流转,尽量采用基于事件的最终一致,降低微服务之间直接访问的压力,实现微服务之间的解耦,维护领域模型的独立性和数据一致性。
除此之外,领域事件驱动机制可以实现一个发布方 N 个订阅方的模式,这在传统的直接服务调用设计中基本是不可能做到的。
CQRS Core Principles — Architecture Integration
核心分离原则
传统 CRUD 与 CQRS 对比:
传统 CRUD: CQRS:
┌──────────┐ ┌──────────┐ ┌──────────┐
│ CRUD │ │ Command │ │ Query │
│ Service │ │ Model │ │ Model │
└────┬─────┘ └────┬─────┘ └────┬─────┘
│ │ │
┌────┴────┐ ┌───┴───┐ ┌───┴───┐
│ DB │ │ Write │ │ Read │
└─────────┘ │ DB │ │ DB │
└───────┘ └───────┘L1 — Model Separation Code Example
// 写侧
@Service
public class OrderCommandService {
private final OrderRepository orderRepository;
@Transactional
public OrderCreatedResult createOrder(CreateOrderCommand command) {
Order order = Order.create(command);
orderRepository.save(order);
return OrderCreatedResult.from(order);
}
}
// 读侧
@Service
public class OrderQueryService {
private final OrderReadRepository orderReadRepository;
public OrderDetailDTO getOrderDetail(String orderId) {
return orderReadRepository.findDetailById(orderId)
.orElseThrow(() -> new OrderNotFoundException(orderId));
}
public Page<OrderSummaryDTO> listOrders(OrderQuery query, Pageable pageable) {
return orderReadRepository.findByCriteria(query, pageable);
}
}L2 — Database Separation Code Example
// 领域事件
public class OrderPaidEvent extends DomainEvent {
private final String orderId;
private final BigDecimal totalAmount;
public OrderPaidEvent(String orderId, BigDecimal totalAmount) {
super(orderId);
this.orderId = orderId;
this.totalAmount = totalAmount;
}
}
// 事件处理器 → 同步到 Query DB
@Component
public class OrderPaidEventHandler {
private final OrderReadRepository readRepo;
@EventListener
@Async
public void on(OrderPaidEvent event) {
OrderDocument doc = OrderDocument.from(event);
readRepo.save(doc); // 写入 Elasticsearch / Query DB
}
}L3 — Event Sourcing Code Example
// 事件溯源聚合
public class Order extends EventSourcedAggregate {
private OrderStatus status;
public void pay() {
apply(new OrderPaidEvent(this.id));
}
@EventHandler
private void on(OrderPaidEvent event) {
this.status = OrderStatus.PAID;
}
}
// 投影
@Component
public class OrderProjection {
private final JdbcTemplate jdbc;
@EventListener
public void on(OrderPaidEvent event) {
jdbc.update(
"UPDATE order_projection SET status = ? WHERE id = ?",
"PAID", event.getOrderId()
);
}
}落地步骤
Phase 1: 评估(1d)
→ 确认是否需要 CQRS → 选择 L1/L2/L3 级别
Phase 2: 命令模型设计(1-2d)
→ 设计 Command 对象 → 实现 CommandHandler → 发布领域事件
Phase 3: 查询模型设计(1-2d)
→ 设计 Query 对象 → 实现 QueryHandler → DTO 组装
Phase 4: 事件同步(L2/L3 需要,1-3d)
→ Outbox 表 → 事件发布器 → 投影更新
→ 幂等策略实现 → 重试/死信机制
Phase 5: 集成测试(1-2d)
→ Command → 领域事件 → Query DB 同步 → 验证最终一致领域事件深入 — 事件驱动设计原则
事件生命周期
领域行为 → 构建 DomainEvent → 持久化 → EventBus 发布
↓
┌────────┴────────┐
↓ ↓
本地处理器 (同步) MQ 外发 (异步)
(同一进程) (跨服务)
↓
外部处理器
(幂等 + 补偿)发布策略对比
| 策略 | 延迟 | 复杂度 | 场景 |
|---|---|---|---|
| 轮询发布(定时扫表) | 1-5s | 低 | 中小流量 |
| CDC(Debezium binlog) | < 100ms | 中 | 高流量低延迟 |
| 事务提交回调 | < 10ms | 低 | Spring 项目 |
幂等设计代码示例
// 策略 1: 事件去重表 — UNIQUE 约束防重
@Transactional
public void handle(OrderPaidEvent event) {
eventRecordDao.insert(new EventRecord(event.getEventId()));
// 业务处理...
}
// 策略 2: 状态机守卫 — 状态不可逆,已处理则跳过
public void pay() {
if (this.status == OrderStatus.PAID) return;
this.status = OrderStatus.PAID;
addDomainEvent(new OrderPaidEvent(this.id));
}领域事件设计原则
1. 事件用过去时命名(OrderPaidEvent, OrderShippedEvent) 2. 事件体包含足够上下文,但不含敏感数据 3. 事件版本向后兼容(仅新增字段,不删除/重命名) 4. 事件消费者必须幂等(at-least-once 投递) 5. 聚合内强一致,聚合间最终一致