本文目录导读:

- 项目结构
- Maven依赖配置 (pom.xml)
- 主应用程序类
- 命令模型 (Command Side)
- 事件模型 (Events)
- 领域模型 (Aggregate)
- 查询模型 (Query Side)
- 仓储层
- 查询处理器
- 事件处理器
- REST控制器
- 配置文件
- 测试
- 使用说明
- 架构说明
我来为您提供一个完整的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模式,包括命令处理、事件驱动和查询优化。