我将为你展示一个完整的CQRS(命令查询职责分离)实际案例,以电商订单系统为例。

电商订单系统CQRS实践案例
业务场景概述
一个电商系统需要处理订单的写入操作(创建订单、更新状态)和查询操作(查看订单列表、订单详情、数据分析等)。
痛点:
- 查询场景复杂(多条件筛选、聚合统计)
- 读写负载不均衡
- 传统CRUD导致查询慢、写操作受影响
系统架构设计
┌─────────────────────────────────────────────────────────────┐
│ 客户端/API网关 │
└──────────────┬──────────────────────────┬───────────────────┘
│ │
┌────────▼─────────┐ ┌─────────▼──────────┐
│ Command Side │ │ Query Side │
│ (写模型) │ │ (读模型) │
└────────┬─────────┘ └─────────┬──────────┘
│ │
┌────────▼─────────┐ ┌─────────▼──────────┐
│ OrderCommand │ │ OrderQuery │
│ Service │ │ Service │
└────────┬─────────┘ └─────────┬──────────┘
│ │
┌────────▼─────────┐ ┌─────────▼──────────┐
│ Write DB │ │ Read DB │
│ (MySQL) │ │ (MongoDB/ES) │
│ 3NF规范化 │ │ 反规范化 │
└────────┬─────────┘ └─────────┬──────────┘
│ │
└────────────┬─────────────┘
│
┌──────▼──────┐
│ Event Bus │
│ (RabbitMQ) │
└─────────────┘
核心代码实现
1 命令侧(写模型)
命令对象定义
// 创建订单命令
@Data
public class CreateOrderCommand {
private String customerId;
private List<OrderItem> items;
private Address shippingAddress;
private String paymentMethod;
}
// 订单项
@Data
public class OrderItem {
private String productId;
private Integer quantity;
private BigDecimal price;
}
命令处理器
@Service
@RequiredArgsConstructor
public class OrderCommandService {
private final OrderRepository orderRepository;
private final ProductServiceClient productServiceClient;
private final EventPublisher eventPublisher;
private final IdGenerator idGenerator;
@Transactional
public String createOrder(CreateOrderCommand command) {
// 1. 验证业务规则
validateOrder(command);
// 2. 创建订单实体(聚合根)
Order order = Order.create(
idGenerator.generate(),
command.getCustomerId(),
command.getItems(),
command.getShippingAddress(),
command.getPaymentMethod()
);
// 3. 计算订单价格(远程调用产品服务)
BigDecimal totalPrice = calculateTotalPrice(command.getItems());
order.calculateTotal(totalPrice);
// 4. 保存订单(写模型DB)
orderRepository.save(order);
// 5. 发布领域事件
eventPublisher.publish(new OrderCreatedEvent(
order.getId(),
order.getCustomerId(),
order.getTotalAmount(),
order.getCreatedAt()
));
return order.getId();
}
@Transactional
public void updateOrderStatus(UpdateOrderStatusCommand command) {
Order order = orderRepository.findById(command.getOrderId());
// 业务规则校验
if (!order.canUpdateTo(command.getNewStatus())) {
throw new InvalidOrderStatusException("无效的状态变更");
}
order.updateStatus(command.getNewStatus());
orderRepository.save(order);
// 发布状态更新事件
eventPublisher.publish(new OrderStatusUpdatedEvent(
order.getId(),
command.getNewStatus(),
command.getOperatorId()
));
}
}
领域模型(聚合根)
@Entity
@Table(name = "orders")
public class Order {
@Id
private String id;
private String customerId;
@Enumerated(EnumType.STRING)
private OrderStatus status;
private BigDecimal totalAmount;
private LocalDateTime createdAt;
private LocalDateTime updatedAt;
@Version
private Long version; // 乐观锁
// 保护方法
protected Order() {}
public static Order create(String id, String customerId,
List<OrderItem> items,
Address shippingAddress,
String paymentMethod) {
Order order = new Order();
order.id = id;
order.customerId = customerId;
order.status = OrderStatus.CREATED;
order.createdAt = LocalDateTime.now();
return order;
}
public void updateStatus(OrderStatus newStatus) {
this.status = newStatus;
this.updatedAt = LocalDateTime.now();
}
public boolean canUpdateTo(OrderStatus newStatus) {
// 状态机校验逻辑
return this.status != OrderStatus.CANCELLED
&& this.status != OrderStatus.COMPLETED;
}
}
2 查询侧(读模型)
查询服务
@Service
@RequiredArgsConstructor
public class OrderQueryService {
private final OrderReadRepository orderReadRepository;
private final OrderSearchRepository searchRepository;
// 获取订单列表
public Page<OrderDTO> getOrders(OrderQuery query) {
// 使用专用的查询数据库/索引
return orderReadRepository.findOrders(
query.getCustomerId(),
query.getStatus(),
query.getStartDate(),
query.getEndDate(),
query.getPageable()
);
}
// 获取订单详情
public OrderDetailDTO getOrderDetail(String orderId) {
return orderReadRepository.findById(orderId)
.orElseThrow(() -> new OrderNotFoundException(orderId));
}
// 复杂聚合查询(使用Elasticsearch)
public List<OrderAnalyticsDTO> getSalesAnalytics(DateRange range) {
return searchRepository.aggregateByDate(range);
}
// 热销商品排行
public List<TopProductDTO> getTopProducts(int limit, DateRange range) {
return searchRepository.getTopProducts(limit, range);
}
}
读模型DTO(反规范化)
// 订单列表DTO - 为列表页优化
@Data
public class OrderListDTO {
private String orderId;
private String customerId;
private String customerName; // 冗余存储
private BigDecimal totalAmount;
private OrderStatus status;
private LocalDateTime createdAt;
private Integer itemCount;
}
读模型Repository(使用MongoDB)
@Repository
public interface OrderReadRepository extends MongoRepository<OrderDocument, String> {
// 组合查询
List<OrderDocument> findByCustomerIdAndStatusAndCreatedAtBetween(
String customerId,
OrderStatus status,
LocalDateTime start,
LocalDateTime end
);
@Aggregation(pipeline = {
"{$group: {_id: '$status', count: {$sum: 1}, totalAmount: {$sum: '$totalAmount'}}}"
})
List<OrderAggregationResult> aggregateByStatus();
}
3 事件处理与数据同步
事件处理器(同步写模型到读模型)
@Component
@RequiredArgsConstructor
public class OrderEventProcessor {
private final OrderReadRepository readRepository;
private final CustomerDetailRepository customerDetailRepository;
@EventListener
@Transactional
public void handleOrderCreated(OrderCreatedEvent event) {
// 将事件数据转换为读模型
OrderDocument orderDoc = new OrderDocument();
orderDoc.setOrderId(event.getOrderId());
orderDoc.setCustomerId(event.getCustomerId());
orderDoc.setTotalAmount(event.getTotalAmount());
orderDoc.setStatus(OrderStatus.CREATED);
orderDoc.setCreatedAt(event.getCreatedAt());
// 冗余客户信息以加速查询
CustomerBrief customer = customerDetailRepository.getBriefCustomer(event.getCustomerId());
orderDoc.setCustomerName(customer.getName());
orderDoc.setCustomerLevel(customer.getLevel());
readRepository.save(orderDoc);
}
@EventListener
@Transactional
public void handleOrderStatusUpdated(OrderStatusUpdatedEvent event) {
// 更新读模型的状态
readRepository.findById(event.getOrderId()).ifPresent(orderDoc -> {
orderDoc.setStatus(event.getNewStatus());
orderDoc.setUpdatedAt(event.getUpdatedAt());
readRepository.save(orderDoc);
});
}
}
4 API控制器
@RestController
@RequestMapping("/api/orders")
@RequiredArgsConstructor
public class OrderController {
private final OrderCommandService commandService;
private final OrderQueryService queryService;
// 命令端点(写操作)
@PostMapping
public ResponseEntity<String> createOrder(@RequestBody CreateOrderCommand command) {
String orderId = commandService.createOrder(command);
return ResponseEntity.status(HttpStatus.CREATED).body(orderId);
}
@PutMapping("/{orderId}/status")
public ResponseEntity<Void> updateStatus(
@PathVariable String orderId,
@RequestBody UpdateOrderStatusCommand command) {
command.setOrderId(orderId);
commandService.updateOrderStatus(command);
return ResponseEntity.ok().build();
}
// 查询端点(读操作)
@GetMapping
public ResponseEntity<Page<OrderListDTO>> getOrders(OrderQuery query) {
return ResponseEntity.ok(queryService.getOrders(query));
}
@GetMapping("/{orderId}")
public ResponseEntity<OrderDetailDTO> getOrderDetail(@PathVariable String orderId) {
return ResponseEntity.ok(queryService.getOrderDetail(orderId));
}
@GetMapping("/analytics")
public ResponseEntity<List<OrderAnalyticsDTO>> getAnalytics(DateRange range) {
return ResponseEntity.ok(queryService.getSalesAnalytics(range));
}
}
数据一致性方案
// 最终一致性实现 - 重试机制
@Component
public class EventReplayService {
private final IntegrationEventLogRepository eventLogRepository;
private final OrderReadRepository readRepository;
@Scheduled(fixedDelay = 60000) // 每分钟检查一次
public void replayFailedEvents() {
List<IntegrationEventLog> failedEvents =
eventLogRepository.findByStatus(EventStatus.FAILED);
for (IntegrationEventLog eventLog : failedEvents) {
try {
// 重新处理事件
processEvent(eventLog.getEventData());
eventLog.markAsProcessed();
} catch (Exception e) {
log.error("Event replay failed: {}", eventLog.getId(), e);
eventLog.incrementRetryCount();
}
}
}
}
性能优化措施
// 查询缓存
@Configuration
public class CacheConfig {
@Bean
public CacheManager cacheManager() {
return new ConcurrentMapCacheManager(
"orderSummaries",
"customerOrders"
);
}
}
@Service
public class OrderQueryService {
@Cacheable(value = "orderSummaries", key = "#orderId")
public OrderDetailDTO getOrderDetailCached(String orderId) {
return orderReadRepository.findById(orderId)
.map(this::convertToDetailDTO)
.orElseThrow(() -> new OrderNotFoundException(orderId));
}
}
监控与治理
// 读写指标监控
@Component
public class CQRSMetricsMonitor {
private final MeterRegistry meterRegistry;
public void recordCommandMetrics(String commandType, long durationMillis) {
meterRegistry.counter(
"order.commands",
"type", commandType
).increment();
meterRegistry.timer(
"order.command.duration",
"type", commandType
).record(Duration.ofMillis(durationMillis));
}
public void recordQueryMetrics(String queryType, long durationMillis) {
meterRegistry.counter(
"order.queries",
"type", queryType
).increment();
meterRegistry.timer(
"order.query.duration",
"type", queryType
).record(Duration.ofMillis(durationMillis));
}
}
优缺点分析
优势:
- ✅ 读写独立扩展,性能优化更容易
- ✅ 读模型可为不同业务场景定制
- ✅ 写模型保持业务逻辑的完整性和一致性
- ✅ 降低数据库的并发压力
挑战:
- ⚠️ 数据最终一致性,需要处理同步延迟
- ⚠️ 系统复杂度增加,需要额外维护
- ⚠️ 事件处理和重放机制需要完善
- ⚠️ 运维成本较高
适用场景建议
适合使用CQRS的场景:
- 读操作远多于写操作的系统
- 复杂的查询和分析需求
- 需要高可用和水平扩展
- 有明确的事件驱动需求
不适合的场景:
- 简单CRUD系统
- 数据实时一致性要求极高
- 团队规模小,资源有限
这个案例展示了CQRS在实际项目中的完整实现,包括命令/查询分离、事件驱动、数据同步和一致性保障等核心要素。