CQRS案例

wen java案例 1

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

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在实际项目中的完整实现,包括命令/查询分离、事件驱动、数据同步和一致性保障等核心要素。

抱歉,评论功能暂时关闭!