Java CQRS模式案例

wen java案例 2

本文目录导读:

Java CQRS模式案例

  1. 项目结构
  2. Maven依赖配置 (pom.xml)
  3. 主应用程序类
  4. 命令模型 (Command Side)
  5. 事件模型 (Events)
  6. 领域模型 (Aggregate)
  7. 查询模型 (Query Side)
  8. 仓储层
  9. 查询处理器
  10. 事件处理器
  11. REST控制器
  12. 配置文件
  13. 测试
  14. 使用说明
  15. 架构说明

我来为您提供一个完整的Java CQRS模式案例,使用Spring Boot和Axon框架实现。

项目结构

cqrs-example/
├── pom.xml
├── src/main/java/com/example/cqrs/
│   ├── CqrsApplication.java
│   ├── command/
│   │   ├── CreateOrderCommand.java
│   │   ├── UpdateOrderStatusCommand.java
│   │   └── CancelOrderCommand.java
│   ├── query/
│   │   ├── GetOrderQuery.java
│   │   └── OrderSummary.java
│   ├── events/
│   │   ├── OrderCreatedEvent.java
│   │   ├── OrderStatusUpdatedEvent.java
│   │   └── OrderCancelledEvent.java
│   ├── model/
│   │   ├── Order.java
│   │   └── OrderStatus.java
│   ├── commandhandlers/
│   │   └── OrderCommandHandler.java
│   ├── queryhandlers/
│   │   └── OrderQueryHandler.java
│   └── repository/
│       └── OrderSummaryRepository.java
└── src/main/resources/
    └── application.yml

Maven依赖配置 (pom.xml)

<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
         xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 
         http://maven.apache.org/xsd/maven-4.0.0.xsd">
    <modelVersion>4.0.0</modelVersion>
    <parent>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-parent</artifactId>
        <version>2.6.3</version>
        <relativePath/>
    </parent>
    <groupId>com.example</groupId>
    <artifactId>cqrs-example</artifactId>
    <version>1.0.0</version>
    <name>cqrs-example</name>
    <description>CQRS Pattern Example with Axon and Spring Boot</description>
    <properties>
        <java.version>11</java.version>
        <axon.version>4.5.5</axon.version>
    </properties>
    <dependencies>
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-web</artifactId>
        </dependency>
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-data-jpa</artifactId>
        </dependency>
        <dependency>
            <groupId>org.axonframework</groupId>
            <artifactId>axon-spring-boot-starter</artifactId>
            <version>${axon.version}</version>
        </dependency>
        <dependency>
            <groupId>org.axonframework</groupId>
            <artifactId>axon-spring-boot-autoconfigure</artifactId>
            <version>${axon.version}</version>
        </dependency>
        <dependency>
            <groupId>com.h2database</groupId>
            <artifactId>h2</artifactId>
            <scope>runtime</scope>
        </dependency>
        <dependency>
            <groupId>org.projectlombok</groupId>
            <artifactId>lombok</artifactId>
            <optional>true</optional>
        </dependency>
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-test</artifactId>
            <scope>test</scope>
        </dependency>
    </dependencies>
</project>

主应用程序类

package com.example.cqrs;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
@SpringBootApplication
public class CqrsApplication {
    public static void main(String[] args) {
        SpringApplication.run(CqrsApplication.class, args);
    }
}

命令模型 (Command Side)

CreateOrderCommand.java

package com.example.cqrs.command;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import org.axonframework.modelling.command.TargetAggregateIdentifier;
@Data
@AllArgsConstructor
@NoArgsConstructor
public class CreateOrderCommand {
    @TargetAggregateIdentifier
    private String orderId;
    private String customerId;
    private String productName;
    private Double price;
    private Integer quantity;
}

UpdateOrderStatusCommand.java

package com.example.cqrs.command;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import org.axonframework.modelling.command.TargetAggregateIdentifier;
@Data
@AllArgsConstructor
@NoArgsConstructor
public class UpdateOrderStatusCommand {
    @TargetAggregateIdentifier
    private String orderId;
    private String newStatus;
}

CancelOrderCommand.java

package com.example.cqrs.command;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import org.axonframework.modelling.command.TargetAggregateIdentifier;
@Data
@AllArgsConstructor
@NoArgsConstructor
public class CancelOrderCommand {
    @TargetAggregateIdentifier
    private String orderId;
    private String reason;
}

事件模型 (Events)

OrderCreatedEvent.java

package com.example.cqrs.events;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
@Data
@AllArgsConstructor
@NoArgsConstructor
public class OrderCreatedEvent {
    private String orderId;
    private String customerId;
    private String productName;
    private Double price;
    private Integer quantity;
    private String status;
}

OrderStatusUpdatedEvent.java

package com.example.cqrs.events;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
@Data
@AllArgsConstructor
@NoArgsConstructor
public class OrderStatusUpdatedEvent {
    private String orderId;
    private String newStatus;
}

OrderCancelledEvent.java

package com.example.cqrs.events;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
@Data
@AllArgsConstructor
@NoArgsConstructor
public class OrderCancelledEvent {
    private String orderId;
    private String reason;
    private String status;
}

领域模型 (Aggregate)

Order.java

package com.example.cqrs.model;
import com.example.cqrs.command.CancelOrderCommand;
import com.example.cqrs.command.CreateOrderCommand;
import com.example.cqrs.command.UpdateOrderStatusCommand;
import com.example.cqrs.events.OrderCancelledEvent;
import com.example.cqrs.events.OrderCreatedEvent;
import com.example.cqrs.events.OrderStatusUpdatedEvent;
import lombok.NoArgsConstructor;
import org.axonframework.commandhandling.CommandHandler;
import org.axonframework.eventsourcing.EventSourcingHandler;
import org.axonframework.modelling.command.AggregateIdentifier;
import org.axonframework.modelling.command.AggregateLifecycle;
import org.axonframework.spring.stereotype.Aggregate;
@Aggregate
@NoArgsConstructor
public class Order {
    @AggregateIdentifier
    private String orderId;
    private String customerId;
    private String productName;
    private Double price;
    private Integer quantity;
    private Double totalAmount;
    private OrderStatus status;
    private String cancelReason;
    @CommandHandler
    public Order(CreateOrderCommand command) {
        // 校验业务规则
        if (command.getQuantity() <= 0) {
            throw new IllegalArgumentException("Quantity must be positive");
        }
        if (command.getPrice() <= 0) {
            throw new IllegalArgumentException("Price must be positive");
        }
        // 触发创建订单事件
        AggregateLifecycle.apply(new OrderCreatedEvent(
            command.getOrderId(),
            command.getCustomerId(),
            command.getProductName(),
            command.getPrice(),
            command.getQuantity(),
            OrderStatus.PENDING.toString()
        ));
    }
    @CommandHandler
    public void handle(UpdateOrderStatusCommand command) {
        if (status == OrderStatus.CANCELLED) {
            throw new IllegalStateException("Cancelled orders cannot be updated");
        }
        AggregateLifecycle.apply(new OrderStatusUpdatedEvent(
            command.getOrderId(),
            command.getNewStatus()
        ));
    }
    @CommandHandler
    public void handle(CancelOrderCommand command) {
        if (status == OrderStatus.CANCELLED) {
            throw new IllegalStateException("Order is already cancelled");
        }
        AggregateLifecycle.apply(new OrderCancelledEvent(
            command.getOrderId(),
            command.getReason(),
            OrderStatus.CANCELLED.toString()
        ));
    }
    @EventSourcingHandler
    public void on(OrderCreatedEvent event) {
        this.orderId = event.getOrderId();
        this.customerId = event.getCustomerId();
        this.productName = event.getProductName();
        this.price = event.getPrice();
        this.quantity = event.getQuantity();
        this.totalAmount = event.getPrice() * event.getQuantity();
        this.status = OrderStatus.valueOf(event.getStatus());
    }
    @EventSourcingHandler
    public void on(OrderStatusUpdatedEvent event) {
        this.status = OrderStatus.valueOf(event.getNewStatus());
    }
    @EventSourcingHandler
    public void on(OrderCancelledEvent event) {
        this.status = OrderStatus.valueOf(event.getStatus());
        this.cancelReason = event.getReason();
    }
}

OrderStatus.java

package com.example.cqrs.model;
public enum OrderStatus {
    PENDING,
    CONFIRMED,
    SHIPPED,
    DELIVERED,
    CANCELLED
}

查询模型 (Query Side)

OrderSummary.java

package com.example.cqrs.query;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import javax.persistence.Entity;
import javax.persistence.Id;
import javax.persistence.Table;
@Data
@AllArgsConstructor
@NoArgsConstructor
@Entity
@Table(name = "order_summary")
public class OrderSummary {
    @Id
    private String orderId;
    private String customerId;
    private String productName;
    private Double price;
    private Integer quantity;
    private Double totalAmount;
    private String status;
    private String cancelReason;
    public OrderSummary(String orderId, String customerId, String productName, 
                        Double price, Integer quantity, Double totalAmount, String status) {
        this.orderId = orderId;
        this.customerId = customerId;
        this.productName = productName;
        this.price = price;
        this.quantity = quantity;
        this.totalAmount = totalAmount;
        this.status = status;
    }
}

GetOrderQuery.java

package com.example.cqrs.query;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
@Data
@AllArgsConstructor
@NoArgsConstructor
public class GetOrderQuery {
    private String orderId;
}

GetAllOrdersQuery.java

package com.example.cqrs.query;
import lombok.Data;
@Data
public class GetAllOrdersQuery {
    private Integer page;
    private Integer size;
    public GetAllOrdersQuery() {
        this.page = 0;
        this.size = 10;
    }
    public GetAllOrdersQuery(Integer page, Integer size) {
        this.page = page != null ? page : 0;
        this.size = size != null ? size : 10;
    }
}

仓储层

OrderSummaryRepository.java

package com.example.cqrs.repository;
import com.example.cqrs.query.OrderSummary;
import org.springframework.data.jpa.repository.JpaRepository;
import org.springframework.stereotype.Repository;
import java.util.List;
@Repository
public interface OrderSummaryRepository extends JpaRepository<OrderSummary, String> {
    List<OrderSummary> findByCustomerId(String customerId);
    List<OrderSummary> findByStatus(String status);
    List<OrderSummary> findTop5ByOrderByTotalAmountDesc();
    Long countByStatus(String status);
}

查询处理器

OrderQueryHandler.java

package com.example.cqrs.queryhandlers;
import com.example.cqrs.query.GetAllOrdersQuery;
import com.example.cqrs.query.GetOrderQuery;
import com.example.cqrs.query.OrderSummary;
import com.example.cqrs.repository.OrderSummaryRepository;
import lombok.RequiredArgsConstructor;
import org.axonframework.queryhandling.QueryHandler;
import org.springframework.data.domain.Page;
import org.springframework.data.domain.PageRequest;
import org.springframework.stereotype.Component;
import java.util.List;
import java.util.Optional;
@Component
@RequiredArgsConstructor
public class OrderQueryHandler {
    private final OrderSummaryRepository orderSummaryRepository;
    @QueryHandler
    public Optional<OrderSummary> handle(GetOrderQuery query) {
        return orderSummaryRepository.findById(query.getOrderId());
    }
    @QueryHandler
    public Page<OrderSummary> handle(GetAllOrdersQuery query) {
        return orderSummaryRepository.findAll(
            PageRequest.of(query.getPage(), query.getSize())
        );
    }
}

事件处理器

OrderEventProcessor.java

package com.example.cqrs.eventhandlers;
import com.example.cqrs.events.OrderCancelledEvent;
import com.example.cqrs.events.OrderCreatedEvent;
import com.example.cqrs.events.OrderStatusUpdatedEvent;
import com.example.cqrs.query.OrderSummary;
import com.example.cqrs.repository.OrderSummaryRepository;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.axonframework.eventhandling.EventHandler;
import org.springframework.stereotype.Component;
import java.util.Optional;
@Slf4j
@Component
@RequiredArgsConstructor
public class OrderEventProcessor {
    private final OrderSummaryRepository orderSummaryRepository;
    @EventHandler
    public void handle(OrderCreatedEvent event) {
        log.info("Processing OrderCreatedEvent for order: {}", event.getOrderId());
        OrderSummary orderSummary = new OrderSummary(
            event.getOrderId(),
            event.getCustomerId(),
            event.getProductName(),
            event.getPrice(),
            event.getQuantity(),
            event.getPrice() * event.getQuantity(),
            event.getStatus()
        );
        orderSummaryRepository.save(orderSummary);
    }
    @EventHandler
    public void handle(OrderStatusUpdatedEvent event) {
        log.info("Processing OrderStatusUpdatedEvent for order: {}", event.getOrderId());
        Optional<OrderSummary> orderSummaryOpt = orderSummaryRepository.findById(event.getOrderId());
        orderSummaryOpt.ifPresent(orderSummary -> {
            orderSummary.setStatus(event.getNewStatus());
            orderSummaryRepository.save(orderSummary);
        });
    }
    @EventHandler
    public void handle(OrderCancelledEvent event) {
        log.info("Processing OrderCancelledEvent for order: {}", event.getOrderId());
        Optional<OrderSummary> orderSummaryOpt = orderSummaryRepository.findById(event.getOrderId());
        orderSummaryOpt.ifPresent(orderSummary -> {
            orderSummary.setStatus(event.getStatus());
            orderSummary.setCancelReason(event.getReason());
            orderSummaryRepository.save(orderSummary);
        });
    }
}

REST控制器

OrderController.java

package com.example.cqrs.controller;
import com.example.cqrs.command.CancelOrderCommand;
import com.example.cqrs.command.CreateOrderCommand;
import com.example.cqrs.command.UpdateOrderStatusCommand;
import com.example.cqrs.query.GetAllOrdersQuery;
import com.example.cqrs.query.GetOrderQuery;
import com.example.cqrs.query.OrderSummary;
import lombok.RequiredArgsConstructor;
import org.axonframework.commandhandling.gateway.CommandGateway;
import org.axonframework.queryhandling.QueryGateway;
import org.springframework.data.domain.Page;
import org.springframework.http.HttpStatus;
import org.springframework.web.bind.annotation.*;
import java.util.UUID;
import java.util.concurrent.CompletableFuture;
@RestController
@RequestMapping("/api/orders")
@RequiredArgsConstructor
public class OrderController {
    private final CommandGateway commandGateway;
    private final QueryGateway queryGateway;
    @PostMapping
    @ResponseStatus(HttpStatus.CREATED)
    public CompletableFuture<String> createOrder(@RequestBody CreateOrderRequest request) {
        String orderId = UUID.randomUUID().toString();
        CreateOrderCommand command = new CreateOrderCommand(
            orderId,
            request.getCustomerId(),
            request.getProductName(),
            request.getPrice(),
            request.getQuantity()
        );
        return commandGateway.send(command)
            .thenApply(result -> orderId);
    }
    @PutMapping("/{orderId}/status")
    public CompletableFuture<Void> updateOrderStatus(@PathVariable String orderId, 
                                                     @RequestBody UpdateStatusRequest request) {
        UpdateOrderStatusCommand command = new UpdateOrderStatusCommand(
            orderId,
            request.getNewStatus()
        );
        return commandGateway.send(command);
    }
    @PostMapping("/{orderId}/cancel")
    public CompletableFuture<Void> cancelOrder(@PathVariable String orderId, 
                                               @RequestBody CancelRequest request) {
        CancelOrderCommand command = new CancelOrderCommand(
            orderId,
            request.getReason()
        );
        return commandGateway.send(command);
    }
    @GetMapping("/{orderId}")
    public CompletableFuture<Optional<OrderSummary>> getOrder(@PathVariable String orderId) {
        GetOrderQuery query = new GetOrderQuery(orderId);
        return queryGateway.query(query, Optional.class);
    }
    @GetMapping
    public CompletableFuture<Page<OrderSummary>> getAllOrders(
            @RequestParam(defaultValue = "0") int page,
            @RequestParam(defaultValue = "10") int size) {
        GetAllOrdersQuery query = new GetAllOrdersQuery(page, size);
        return queryGateway.query(query, Page.class);
    }
    // 内部类用于请求体
    static class CreateOrderRequest {
        public String customerId;
        public String productName;
        public Double price;
        public Integer quantity;
    }
    static class UpdateStatusRequest {
        public String newStatus;
    }
    static class CancelRequest {
        public String reason;
    }
}

配置文件

application.yml

server:
  port: 8080
spring:
  datasource:
    url: jdbc:h2:mem:cqrsdb
    driverClassName: org.h2.Driver
    username: sa
    password: 
  jpa:
    hibernate:
      ddl-auto: create-drop
    show-sql: true
  h2:
    console:
      enabled: true
      path: /h2-console
axon:
  eventhandling:
    processors:
      order-group:
        mode: SUBSCRIBING
  serializer:
    general: jackson
    events: jackson
logging:
  level:
    com.example.cqrs: DEBUG
    org.axonframework: INFO

测试

OrderControllerTest.java

package com.example.cqrs;
import com.example.cqrs.controller.OrderController;
import com.example.cqrs.query.OrderSummary;
import com.example.cqrs.repository.OrderSummaryRepository;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.boot.test.web.client.TestRestTemplate;
import org.springframework.http.HttpEntity;
import org.springframework.http.HttpHeaders;
import org.springframework.http.MediaType;
import org.springframework.http.ResponseEntity;
import java.util.Map;
import java.util.Optional;
import static org.junit.jupiter.api.Assertions.*;
@SpringBootTest(webEnvironment = SpringBootTest.WebEnvironment.RANDOM_PORT)
class OrderControllerTest {
    @Autowired
    private TestRestTemplate restTemplate;
    @Autowired
    private OrderSummaryRepository orderSummaryRepository;
    @Test
    void testCreateAndGetOrder() throws Exception {
        // 创建订单
        Map<String, Object> request = Map.of(
            "customerId", "cust-001",
            "productName", "iPhone 15",
            "price", 999.99,
            "quantity", 2
        );
        ResponseEntity<String> createResponse = restTemplate.postForEntity(
            "/api/orders",
            new HttpEntity<>(request, getJsonHeaders()),
            String.class
        );
        assertEquals(201, createResponse.getStatusCodeValue());
        assertNotNull(createResponse.getBody());
        // 等待事件处理完成
        Thread.sleep(1000);
        // 查询订单
        ResponseEntity<OrderSummary> getResponse = restTemplate.getForEntity(
            "/api/orders/" + createResponse.getBody(),
            OrderSummary.class
        );
        assertTrue(getResponse.getBody() != null);
        assertEquals("cust-001", getResponse.getBody().getCustomerId());
        assertEquals("iPhone 15", getResponse.getBody().getProductName());
        assertEquals(Double.valueOf(1999.98), getResponse.getBody().getTotalAmount());
    }
    @Test
    void testOrderStatusUpdate() throws Exception {
        // 创建订单
        Map<String, Object> createRequest = Map.of(
            "customerId", "cust-002",
            "productName", "MacBook Pro",
            "price", 1299.00,
            "quantity", 1
        );
        ResponseEntity<String> createResponse = restTemplate.postForEntity(
            "/api/orders",
            new HttpEntity<>(createRequest, getJsonHeaders()),
            String.class
        );
        String orderId = createResponse.getBody();
        assertNotNull(orderId);
        // 更新状态
        Map<String, String> statusRequest = Map.of("newStatus", "CONFIRMED");
        restTemplate.put(
            "/api/orders/" + orderId + "/status",
            new HttpEntity<>(statusRequest, getJsonHeaders())
        );
        // 等待事件处理完成
        Thread.sleep(1000);
        // 验证状态更新
        Optional<OrderSummary> orderOpt = orderSummaryRepository.findById(orderId);
        assertTrue(orderOpt.isPresent());
        assertEquals("CONFIRMED", orderOpt.get().getStatus());
    }
    @Test
    void testCancelOrder() throws Exception {
        // 创建订单
        Map<String, Object> createRequest = Map.of(
            "customerId", "cust-003",
            "productName", "iPad",
            "price", 499.00,
            "quantity", 1
        );
        ResponseEntity<String> createResponse = restTemplate.postForEntity(
            "/api/orders",
            new HttpEntity<>(createRequest, getJsonHeaders()),
            String.class
        );
        String orderId = createResponse.getBody();
        assertNotNull(orderId);
        // 取消订单
        Map<String, String> cancelRequest = Map.of("reason", "Customer changed mind");
        restTemplate.postForEntity(
            "/api/orders/" + orderId + "/cancel",
            new HttpEntity<>(cancelRequest, getJsonHeaders()),
            Void.class
        );
        // 等待事件处理完成
        Thread.sleep(1000);
        // 验证取消状态
        Optional<OrderSummary> orderOpt = orderSummaryRepository.findById(orderId);
        assertTrue(orderOpt.isPresent());
        assertEquals("CANCELLED", orderOpt.get().getStatus());
        assertEquals("Customer changed mind", orderOpt.get().getCancelReason());
    }
    private HttpHeaders getJsonHeaders() {
        HttpHeaders headers = new HttpHeaders();
        headers.setContentType(MediaType.APPLICATION_JSON);
        return headers;
    }
}

使用说明

运行应用程序

mvn spring-boot:run

测试API端点

创建订单

curl -X POST http://localhost:8080/api/orders \
  -H "Content-Type: application/json" \
  -d '{
    "customerId": "cust-001",
    "productName": "iPhone 15",
    "price": 999.99,
    "quantity": 2
  }'

获取订单

curl http://localhost:8080/api/orders/{orderId}

更新订单状态

curl -X PUT http://localhost:8080/api/orders/{orderId}/status \
  -H "Content-Type: application/json" \
  -d '{"newStatus": "CONFIRMED"}'

取消订单

curl -X POST http://localhost:8080/api/orders/{orderId}/cancel \
  -H "Content-Type: application/json" \
  -d '{"reason": "Customer changed mind"}'

获取所有订单(分页)

curl "http://localhost:8080/api/orders?page=0&size=10"

访问H2控制台

  • URL: http://localhost:8080/h2-console
  • JDBC URL: jdbc:h2:mem:cqrsdb
  • Username: sa
  • Password: (留空)

架构说明

命令端(写模型)

  • 命令处理器:验证业务规则,执行命令
  • 聚合:订单实体,使用事件溯源
  • 命令网关:发送命令到指定的命令处理器

查询端(读模型)

  • 查询处理器:处理查询请求
  • 查询数据库:使用JPA存储查询模型
  • 查询网关:路由查询请求

事件驱动

  • 事件处理器:更新查询模型
  • 事件总线:在命令端和查询端之间传递事件

核心优势

  • 职责分离:读写操作独立优化
  • 独立扩展:读写数据库可以独立扩展
  • 事件溯源:完整的审计历史
  • 异步处理:提高系统响应性

这个案例展示了如何使用Axon框架在Spring Boot中实现完整的CQRS模式,包括命令处理、事件驱动和查询优化。

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