Triển Khai Kiến Trúc CQRS Trong Spring Boot 3 Với Spring Data JPA và Redis
Kiến trúc CQRS (Command Query Responsibility Segregation) đang ngày càng trở thành lựa chọn hàng đầu cho các hệ thống enterprise có nhu cầu đọc/ghi dữ liệu chênh lệch lớn. Trong bài viết này, chúng ta sẽ cùng nhau xây dựng một Lightweight CQRS hoàn chỉnh trong Spring Boot 3 mà không cần dùng đến các framework cồng kềnh như Axon.
1. CQRS Pattern Là Gì? Khi Nào Nên Và Không Nên Sử Dụng?
Nói một cách dễ hiểu: CQRS tách biệt mô hình dữ liệu dùng để ghi (Command) và mô hình dữ liệu dùng để đọc (Query). Thay vì ép một mô hình dữ liệu duy nhất phải phục vụ cả hai mục đích, chúng ta xây dựng hai mô hình riêng biệt, mỗi mô hình được tối ưu hóa cho nhiệm vụ cụ thể của nó.
Cách tiếp cận truyền thống CRUD với một model duy nhất hoạt động tốt khi tỷ lệ đọc/ghi cân bằng. Tuy nhiên, khi nhu cầu đọc và ghi bắt đầu phân kỳ (divergence), mô hình đơn nhất bắt đầu bộc lộ những hạn chế rõ rệt.
1.1 Khái niệm Command (Write) vs Query (Read)
Trong kiến trúc CQRS:
| Khía cạnh | Command (Write Side) | Query (Read Side) |
|---|---|---|
| Mục đích | Thay đổi trạng thái hệ thống | Truy xuất dữ liệu |
| Model | Domain Model (chuẩn hóa, giàu logic nghiệp vụ) | Read Model (phi chuẩn hóa, tối ưu cho hiển thị) |
| Database | PostgreSQL (transactional, ACID) | Redis (in-memory, siêu nhanh) |
| Phương thức | @Transactional |
@Cacheable / read-only |
| Tần suất | Thấp | Cao |
📌 Lưu ý: CQRS không bắt buộc phải đi kèm với Event Sourcing. Bạn hoàn toàn có thể triển khai CQRS mà không cần Event Sourcing. Trong bài viết này, chúng ta chỉ tập trung vào CQRS, sử dụng Domain Events để đồng bộ giữa hai model một cách đơn giản.
1.2 Bài toán Eventual Consistency (Đồng bộ dữ liệu bất đồng bộ)
Một điểm cần lưu ý (Verification Risk): CQRS thường dẫn đến tính nhất quán cuối cùng (Eventual Consistency). Điều này có nghĩa là sau khi ghi dữ liệu thành công, có thể mất vài mili giây đến vài giây để dữ liệu được đồng bộ sang Read Model. Đây là một sự đánh đổi có chủ đích để đạt được hiệu năng và khả năng mở rộng.
Khi nào bạn không nên sử dụng CQRS?
- Ứng dụng CRUD đơn giản, tỷ lệ đọc/ghi cân bằng
- Yêu cầu nhất quán tức thời (strong consistency) tuyệt đối
- Team chưa có kinh nghiệm với kiến trúc phức tạp
- Chi phí vận hành (maintenance) vượt quá lợi ích mang lại
1.3 So sánh Custom Spring CQRS vs Axon Framework
Điểm khác biệt lớn nhất của bài viết này là chúng ta không sử dụng Axon Framework – một giải pháp phổ biến nhưng khá nặng nề cho CQRS/Event Sourcing. Dưới đây là bảng so sánh chi tiết:
| Tiêu chí | Custom Spring CQRS | Axon Framework |
|---|---|---|
| Độ phức tạp | Thấp, dễ hiểu, ít abstraction | Cao, nhiều tầng abstraction |
| Learning curve | Thấp – chỉ cần biết Spring Boot | Cao – cần học các khái niệm riêng của Axon |
| Tính linh hoạt | Cao – tự do thiết kế theo nhu cầu | Trung bình – bị ràng buộc bởi framework |
| Event Sourcing | Tự implement (nếu cần) | Tích hợp sẵn, mạnh mẽ |
| Distributed / Saga | Cần tự xây dựng | Hỗ trợ sẵn |
| Phù hợp | Ứng dụng đơn service, POC, dự án vừa | Hệ thống phân tán lớn, nhiều service |
| Thời gian startup | Nhanh | Chậm hơn do nhiều component |
Khi nào nên chọn Custom:
- Bạn muốn kiểm soát hoàn toàn kiến trúc
- Dự án không yêu cầu distributed transaction phức tạp
- Team đã quen với Spring Boot, không muốn học thêm framework mới
Khi nào nên chọn Axon:
- Hệ thống microservices phức tạp, cần Event Sourcing và Saga
- Cần các tính năng sẵn có như command bus, event bus, snapshot
- Có nguồn lực để đầu tư học framework
2. Thiết Kế Kiến Trúc Lightweight CQRS Trong Spring Boot 3
Hãy hình dung luồng xử lý trong kiến trúc CQRS mà chúng ta sẽ xây dựng:
┌─────────────┐ ┌──────────────────────┐ ┌─────────────────┐
│ Client │────▶│ OrderController │────▶│ Command Bus │
│ (REST API) │ │ (API Layer) │ │ (In-memory) │
└─────────────┘ └──────────────────────┘ └────────┬────────┘
│
▼
┌─────────────┐ ┌──────────────────────┐ ┌─────────────────┐
│ Client │────▶│ OrderController │────▶│ Query Bus │
│ (REST API) │ │ (API Layer) │ │ (In-memory) │
└─────────────┘ └──────────────────────┘ └────────┬────────┘
│
▼
┌──────────────────────────────────────────────────────────────────┐
│ │
│ ┌─────────────────────┐ ┌──────────────────────────┐ │
│ │ COMMAND SIDE │ │ QUERY SIDE │ │
│ │ (PostgreSQL) │ │ (Redis) │ │
│ │ │ │ │ │
│ │ Command Handler │ │ Query Handler │ │
│ │ + JPA Repository │ │ + Redis Repository │ │
│ │ + Transactional │ │ + @Cacheable │ │
│ └──────────┬──────────┘ └──────────────────────────┘ │
│ │ │
│ ▼ │
│ ┌─────────────────────┐ │
│ │ ApplicationEvent │ │
│ │ Publisher │ │
│ │ (Spring Events) │ │
│ └──────────┬──────────┘ │
│ │ │
│ ▼ │
│ ┌─────────────────────┐ │
│ │ Event Listener │───▶ Cập nhật Redis Read Model │
│ │ (@EventListener) │ │
│ └─────────────────────┘ │
│ │
└──────────────────────────────────────────────────────────────────┘
Các thành phần chính:
- Command Bus: Điều phối Command đến Handler tương ứng
- Command Handler: Xử lý logic nghiệp vụ, lưu vào PostgreSQL, publish Domain Event
- Event Publisher: Spring
ApplicationEventPublisher– nhẹ, không cần message broker - Event Listener: Lắng nghe sự kiện, cập nhật Read Model trong Redis
- Query Bus: Điều phối Query đến Handler tương ứng
- Query Handler: Truy xuất dữ liệu từ Redis với tốc độ cao

3. Chuẩn bị các file cho thực chiến
3.1 Dependencies (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
https://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>3.2.0</version>
<relativePath/>
</parent>
<groupId>com.example</groupId>
<artifactId>cqrs-order-service</artifactId>
<version>1.0.0</version>
<name>cqrs-order-service</name>
<properties>
<java.version>21</java.version>
</properties>
<dependencies>
<!-- Spring Boot Starters -->
<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.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-redis</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-validation</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-actuator</artifactId>
</dependency>
<!-- Database -->
<dependency>
<groupId>org.postgresql</groupId>
<artifactId>postgresql</artifactId>
<scope>runtime</scope>
</dependency>
<dependency>
<groupId>org.apache.commons</groupId>
<artifactId>commons-pool2</artifactId>
</dependency>
<!-- Utilities -->
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<optional>true</optional>
</dependency>
<!-- Test -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-maven-plugin</artifactId>
</plugin>
</plugins>
</build>
</project>
3.2 Cấu hình application.yml
spring:
datasource:
url: jdbc:postgresql://localhost:5432/orderdb
username: postgres
password: postgres
driver-class-name: org.postgresql.Driver
jpa:
hibernate:
ddl-auto: update
show-sql: true
properties:
hibernate:
dialect: org.hibernate.dialect.PostgreSQLDialect
format_sql: true
data:
redis:
host: localhost
port: 6379
timeout: 2000ms
lettuce:
pool:
max-active: 8
max-idle: 8
min-idle: 0
cache:
type: redis
redis:
time-to-live: 3600000 # 1 hour
logging:
level:
com.example.order: DEBUG
org.springframework.data.redis: DEBUG
3.3 Redis Configuration
// infrastructure/config/RedisConfig.java
package com.example.order.infrastructure.config;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.datatype.jsr310.JavaTimeModule;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.data.redis.connection.RedisConnectionFactory;
import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.data.redis.serializer.Jackson2JsonRedisSerializer;
import org.springframework.data.redis.serializer.StringRedisSerializer;
@Configuration
public class RedisConfig {
@Bean
public RedisTemplate<String, Object> redisTemplate(RedisConnectionFactory connectionFactory) {
RedisTemplate<String, Object> template = new RedisTemplate<>();
template.setConnectionFactory(connectionFactory);
// Use Jackson2JsonRedisSerializer for value serialization
ObjectMapper objectMapper = new ObjectMapper();
objectMapper.registerModule(new JavaTimeModule());
Jackson2JsonRedisSerializer
<Object> serializer =
new Jackson2JsonRedisSerializer<>(objectMapper, Object.class);
template.setKeySerializer(new StringRedisSerializer());
template.setValueSerializer(serializer);
template.setHashKeySerializer(new StringRedisSerializer());
template.setHashValueSerializer(serializer);
template.afterPropertiesSet();
return template;
}
}
4. Thực Chiến: Xây Dựng Luồng Đặt Hàng (Order Service)
4.1 Cấu trúc Project
src/main/java/com/example/order/
├── OrderApplication.java (Spring Boot main, có @EnableAsync)
├── api/
│ └── OrderController.java
├── command/
│ ├── CreateOrderCommand.java (Java Record - Immutable)
│ ├── OrderCommandHandler.java
│ └── OrderWriteRepository.java
├── query/
│ ├── GetOrderQuery.java (Java Record - Immutable)
│ ├── OrderQueryHandler.java
│ └── OrderReadRepository.java (Spring Data Redis)
├── domain/
│ ├── Order.java (JPA Entity)
│ └── OrderStatus.java
├── event/
│ ├── OrderCreatedEvent.java (Java Record)
│ └── OrderEventListener.java
├── dto/
│ ├── CreateOrderRequest.java
│ └── OrderResponseDto.java
├── exception/
│ └── OrderNotFoundException.java
└── infrastructure/
├── bus/
│ ├── CommandBus.java
│ ├── CommandHandler.java
│ ├── QueryBus.java
│ └── QueryHandler.java
└── config/
├── RedisConfig.java
└── AsyncConfig.java
4.2 Application Class với @EnableAsync
// OrderApplication.java
package com.example.order;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.scheduling.annotation.EnableAsync;
@SpringBootApplication
@EnableAsync // Bắt buộc để sử dụng @Async
public class OrderApplication {
public static void main(String[] args) {
SpringApplication.run(OrderApplication.class, args);
}
}
4.3 Domain Entity (JPA)
// domain/Order.java
package com.example.order.domain;
import jakarta.persistence.*;
import java.math.BigDecimal;
import java.time.LocalDateTime;
import java.util.UUID;
@Entity
@Table(name = "orders")
public class Order {
@Id
@GeneratedValue(strategy = GenerationType.UUID)
private UUID id;
@Column(nullable = false)
private String customerId;
@Column(nullable = false)
private String productId;
@Column(nullable = false)
private Integer quantity;
@Column(nullable = false, precision = 19, scale = 2)
private BigDecimal totalPrice;
@Enumerated(EnumType.STRING)
@Column(nullable = false)
private OrderStatus status;
@Column(nullable = false)
private LocalDateTime createdAt;
@Column(nullable = false)
private LocalDateTime updatedAt;
protected Order() {} // JPA required
// Factory method
public static Order create(String customerId, String productId,
Integer quantity, BigDecimal unitPrice) {
if (quantity == null || quantity <= 0) {
throw new IllegalArgumentException("Quantity must be positive");
}
if (unitPrice == null || unitPrice.signum() <= 0) {
throw new IllegalArgumentException("Unit price must be positive");
}
Order order = new Order();
order.id = UUID.randomUUID();
order.customerId = customerId;
order.productId = productId;
order.quantity = quantity;
order.totalPrice = unitPrice.multiply(BigDecimal.valueOf(quantity));
order.status = OrderStatus.PENDING;
order.createdAt = LocalDateTime.now();
order.updatedAt = order.createdAt;
return order;
}
// Getters
public UUID getId() { return id; }
public String getCustomerId() { return customerId; }
public String getProductId() { return productId; }
public Integer getQuantity() { return quantity; }
public BigDecimal getTotalPrice() { return totalPrice; }
public OrderStatus getStatus() { return status; }
public LocalDateTime getCreatedAt() { return createdAt; }
public LocalDateTime getUpdatedAt() { return updatedAt; }
}
// domain/OrderStatus.java
package com.example.order.domain;
public enum OrderStatus {
PENDING, CONFIRMED, SHIPPED, DELIVERED, CANCELLED
}
4.4 Command Object (Java 21 Record)
// command/CreateOrderCommand.java
package com.example.order.command;
import java.math.BigDecimal;
public record CreateOrderCommand(
String customerId,
String productId,
Integer quantity,
BigDecimal unitPrice
) {
public CreateOrderCommand {
if (customerId == null || customerId.isBlank()) {
throw new IllegalArgumentException("customerId is required");
}
if (productId == null || productId.isBlank()) {
throw new IllegalArgumentException("productId is required");
}
if (quantity == null || quantity <= 0) {
throw new IllegalArgumentException("quantity must be positive");
}
if (unitPrice == null || unitPrice.signum() <= 0) {
throw new IllegalArgumentException("unitPrice must be positive");
}
}
}
💡 Best Practice: Sử dụng Java Records cho Command và Query DTOs để đảm bảo tính immutable và giảm boilerplate code.
4.5 Write Repository
// command/OrderWriteRepository.java
package com.example.order.command;
import com.example.order.domain.Order;
import org.springframework.data.jpa.repository.JpaRepository;
import org.springframework.stereotype.Repository;
import java.util.UUID;
@Repository
public interface OrderWriteRepository extends JpaRepository<Order, UUID> {
}
4.6 Command Handler
// command/OrderCommandHandler.java
package com.example.order.command;
import com.example.order.domain.Order;
import com.example.order.event.OrderCreatedEvent;
import com.example.order.infrastructure.bus.CommandHandler;
import org.springframework.context.ApplicationEventPublisher;
import org.springframework.stereotype.Component;
import org.springframework.transaction.annotation.Transactional;
import java.util.UUID;
@Component
public class OrderCommandHandler implements CommandHandler<CreateOrderCommand, UUID> {
private final OrderWriteRepository writeRepository;
private final ApplicationEventPublisher eventPublisher;
public OrderCommandHandler(OrderWriteRepository writeRepository,
ApplicationEventPublisher eventPublisher) {
this.writeRepository = writeRepository;
this.eventPublisher = eventPublisher;
}
@Override
@Transactional
public UUID handle(CreateOrderCommand command) {
// 1. Create domain entity using factory method
Order order = Order.create(
command.customerId(),
command.productId(),
command.quantity(),
command.unitPrice()
);
// 2. Persist to PostgreSQL
Order savedOrder = writeRepository.save(order);
// 3. Publish domain event (synchronous within transaction)
OrderCreatedEvent event = new OrderCreatedEvent(
savedOrder.getId(),
savedOrder.getCustomerId(),
savedOrder.getProductId(),
savedOrder.getQuantity(),
savedOrder.getTotalPrice(),
savedOrder.getStatus(),
savedOrder.getCreatedAt()
);
eventPublisher.publishEvent(event);
// 4. Return the order ID
return savedOrder.getId();
}
}

Giải thích code:
@Transactionalđảm bảo toàn bộ operation được bọc trong một transactionApplicationEventPublisherlà cơ chế event nhẹ của Spring, không yêu cầu thêm dependencies- Event được publish trong cùng transaction – nếu transaction rollback, event sẽ không được gửi
4.7 Domain Event
// event/OrderCreatedEvent.java
package com.example.order.event;
import com.example.order.domain.OrderStatus;
import java.math.BigDecimal;
import java.time.LocalDateTime;
import java.util.UUID;
public record OrderCreatedEvent(
UUID orderId,
String customerId,
String productId,
Integer quantity,
BigDecimal totalPrice,
OrderStatus status,
LocalDateTime createdAt
) {
}
4.8 Read Model (Redis)
// query/OrderView.java
package com.example.order.query;
import com.example.order.event.OrderCreatedEvent;
import org.springframework.data.annotation.Id;
import org.springframework.data.redis.core.RedisHash;
import java.math.BigDecimal;
import java.time.LocalDateTime;
@RedisHash(value = "orders", timeToLive = 3600) // TTL 1 hour
public class OrderView {
@Id
private String id;
private String customerId;
private String productId;
private Integer quantity;
private BigDecimal totalPrice;
private String status;
private String createdAt;
private String lastUpdatedAt;
public static OrderView fromEvent(OrderCreatedEvent event) {
OrderView view = new OrderView();
view.id = event.orderId().toString();
view.customerId = event.customerId();
view.productId = event.productId();
view.quantity = event.quantity();
view.totalPrice = event.totalPrice();
view.status = event.status().name();
view.createdAt = event.createdAt().toString();
view.lastUpdatedAt = LocalDateTime.now().toString();
return view;
}
// Getters and setters
public String getId() { return id; }
public void setId(String id) { this.id = id; }
public String getCustomerId() { return customerId; }
public void setCustomerId(String customerId) { this.customerId = customerId; }
public String getProductId() { return productId; }
public void setProductId(String productId) { this.productId = productId; }
public Integer getQuantity() { return quantity; }
public void setQuantity(Integer quantity) { this.quantity = quantity; }
public BigDecimal getTotalPrice() { return totalPrice; }
public void setTotalPrice(BigDecimal totalPrice) { this.totalPrice = totalPrice; }
public String getStatus() { return status; }
public void setStatus(String status) { this.status = status; }
public String getCreatedAt() { return createdAt; }
public void setCreatedAt(String createdAt) { this.createdAt = createdAt; }
public String getLastUpdatedAt() { return lastUpdatedAt; }
public void setLastUpdatedAt(String lastUpdatedAt) { this.lastUpdatedAt = lastUpdatedAt; }
}
4.9 Read Repository (Spring Data Redis)
// query/OrderReadRepository.java
package com.example.order.query;
import org.springframework.data.repository.CrudRepository;
import org.springframework.stereotype.Repository;
import java.util.List;
@Repository
public interface OrderReadRepository extends CrudRepository<OrderView, String> {
List
<OrderView> findByCustomerId(String customerId);
List
<OrderView> findByStatus(String status);
}
4.10 Event Listener với Retry & Idempotent
// event/OrderEventListener.java
package com.example.order.event;
import com.example.order.query.OrderReadRepository;
import com.example.order.query.OrderView;
import lombok.extern.slf4j.Slf4j;
import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.scheduling.annotation.Async;
import org.springframework.stereotype.Component;
import org.springframework.transaction.event.TransactionPhase;
import org.springframework.transaction.event.TransactionalEventListener;
import java.time.Duration;
@Component
@Slf4j
public class OrderEventListener {
private final OrderReadRepository readRepository;
private final RedisTemplate<String, Object> redisTemplate;
public OrderEventListener(OrderReadRepository readRepository,
RedisTemplate<String, Object> redisTemplate) {
this.readRepository = readRepository;
this.redisTemplate = redisTemplate;
}
@Async
@TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT)
public void handleOrderCreated(OrderCreatedEvent event) {
try {
// Idempotent check: avoid duplicate processing
String idempotentKey = "processed:order:" + event.orderId();
Boolean processed = redisTemplate.opsForValue()
.setIfAbsent(idempotentKey, "true", Duration.ofMinutes(5));
if (Boolean.FALSE.equals(processed)) {
log.warn("Order event already processed: {}", event.orderId());
return;
}
// Build and save read model
OrderView orderView = OrderView.fromEvent(event);
readRepository.save(orderView);
log.info("Read model updated successfully for order: {}", event.orderId());
} catch (Exception e) {
log.error("Failed to update read model for order: {}", event.orderId(), e);
// In production, you would send to a Dead Letter Queue (e.g., Kafka DLQ)
// or schedule retry with exponential backoff.
// For simplicity, we just log and let monitoring alert.
}
}
}
💡 Lưu ý: Sử dụng @TransactionalEventListener với TransactionPhase.AFTER_COMMIT đảm bảo event chỉ được xử lý sau khi transaction chính đã commit thành công.
4.11 Query Handler
// query/GetOrderQuery.java
package com.example.order.query;
import java.util.UUID;
public record GetOrderQuery(UUID orderId) {
public GetOrderQuery {
if (orderId == null) {
throw new IllegalArgumentException("orderId is required");
}
}
}
// query/OrderQueryHandler.java
package com.example.order.query;
import com.example.order.dto.OrderResponseDto;
import com.example.order.exception.OrderNotFoundException;
import com.example.order.infrastructure.bus.QueryHandler;
import org.springframework.cache.annotation.Cacheable;
import org.springframework.stereotype.Component;
@Component
public class OrderQueryHandler implements QueryHandler<GetOrderQuery, OrderResponseDto> {
private final OrderReadRepository readRepository;
public OrderQueryHandler(OrderReadRepository readRepository) {
this.readRepository = readRepository;
}
@Override
@Cacheable(value = "orders", key = "#query.orderId()")
public OrderResponseDto handle(GetOrderQuery query) {
String id = query.orderId().toString();
OrderView orderView = readRepository.findById(id)
.orElseThrow(() -> new OrderNotFoundException("Order not found: " + id));
return OrderResponseDto.from(orderView);
}
}

💡 Lưu ý: Query Handler truy xuất dữ liệu từ Redis mà không chạm vào PostgreSQL, giúp response time dưới 5ms ngay cả khi traffic lớn.
4.12 CommandBus & QueryBus Implementation
// infrastructure/bus/CommandHandler.java
package com.example.order.infrastructure.bus;
public interface CommandHandler<C, R> {
R handle(C command);
}
// infrastructure/bus/CommandBus.java
package com.example.order.infrastructure.bus;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import java.lang.reflect.ParameterizedType;
import java.lang.reflect.Type;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
@Component
public class CommandBus {
private final Map<Class<?>, CommandHandler<?, ?>> handlerMap = new HashMap<>();
@Autowired
@SuppressWarnings("unchecked")
public void registerHandlers(List<CommandHandler<?, ?>> handlers) {
for (CommandHandler<?, ?> handler : handlers) {
Type[] interfaces = handler.getClass().getGenericInterfaces();
for (Type type : interfaces) {
if (type instanceof ParameterizedType paramType) {
Type rawType = paramType.getRawType();
if (rawType == CommandHandler.class) {
Type commandType = paramType.getActualTypeArguments()[0];
handlerMap.put((Class<?>) commandType, handler);
break;
}
}
}
}
}
@SuppressWarnings("unchecked")
public <C, R> R execute(C command) {
CommandHandler<C, R> handler = (CommandHandler<C, R>) handlerMap.get(command.getClass());
if (handler == null) {
throw new IllegalStateException("No handler registered for command: " + command.getClass());
}
return handler.handle(command);
}
}
// infrastructure/bus/QueryHandler.java
package com.example.order.infrastructure.bus;
public interface QueryHandler<Q, R> {
R handle(Q query);
}
// infrastructure/bus/QueryBus.java
package com.example.order.infrastructure.bus;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import java.lang.reflect.ParameterizedType;
import java.lang.reflect.Type;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
@Component
public class QueryBus {
private final Map<Class<?>, QueryHandler<?, ?>> handlerMap = new HashMap<>();
@Autowired
@SuppressWarnings("unchecked")
public void registerHandlers(List<QueryHandler<?, ?>> handlers) {
for (QueryHandler<?, ?> handler : handlers) {
Type[] interfaces = handler.getClass().getGenericInterfaces();
for (Type type : interfaces) {
if (type instanceof ParameterizedType paramType) {
Type rawType = paramType.getRawType();
if (rawType == QueryHandler.class) {
Type queryType = paramType.getActualTypeArguments()[0];
handlerMap.put((Class<?>) queryType, handler);
break;
}
}
}
}
}
@SuppressWarnings("unchecked")
public <Q, R> R execute(Q query) {
QueryHandler<Q, R> handler = (QueryHandler<Q, R>) handlerMap.get(query.getClass());
if (handler == null) {
throw new IllegalStateException("No handler registered for query: " + query.getClass());
}
return handler.handle(query);
}
}
4.13 DTOs & Exception
// dto/CreateOrderRequest.java
package com.example.order.dto;
import jakarta.validation.constraints.DecimalMin;
import jakarta.validation.constraints.Min;
import jakarta.validation.constraints.NotBlank;
import java.math.BigDecimal;
public record CreateOrderRequest(
@NotBlank(message = "customerId is required")
String customerId,
@NotBlank(message = "productId is required")
String productId,
@Min(value = 1, message = "quantity must be at least 1")
Integer quantity,
@DecimalMin(value = "0.01", message = "unitPrice must be positive")
BigDecimal unitPrice
) {
}
// dto/OrderResponseDto.java
package com.example.order.dto;
import com.example.order.query.OrderView;
import java.math.BigDecimal;
import java.util.UUID;
public record OrderResponseDto(
UUID orderId,
String customerId,
String productId,
Integer quantity,
BigDecimal totalPrice,
String status,
String createdAt
) {
public static OrderResponseDto from(OrderView view) {
return new OrderResponseDto(
UUID.fromString(view.getId()),
view.getCustomerId(),
view.getProductId(),
view.getQuantity(),
view.getTotalPrice(),
view.getStatus(),
view.getCreatedAt()
);
}
}
// exception/OrderNotFoundException.java
package com.example.order.exception;
public class OrderNotFoundException extends RuntimeException {
public OrderNotFoundException(String message) {
super(message);
}
}
4.14 Controller API
// api/OrderController.java
package com.example.order.api;
import com.example.order.command.CreateOrderCommand;
import com.example.order.dto.CreateOrderRequest;
import com.example.order.dto.OrderResponseDto;
import com.example.order.infrastructure.bus.CommandBus;
import com.example.order.infrastructure.bus.QueryBus;
import com.example.order.query.GetOrderQuery;
import jakarta.validation.Valid;
import org.springframework.http.HttpStatus;
import org.springframework.http.ResponseEntity;
import org.springframework.web.bind.annotation.*;
import java.util.UUID;
@RestController
@RequestMapping("/api/orders")
public class OrderController {
private final CommandBus commandBus;
private final QueryBus queryBus;
public OrderController(CommandBus commandBus, QueryBus queryBus) {
this.commandBus = commandBus;
this.queryBus = queryBus;
}
@PostMapping
public ResponseEntity
<UUID> createOrder(@Valid @RequestBody CreateOrderRequest request) {
CreateOrderCommand command = new CreateOrderCommand(
request.customerId(),
request.productId(),
request.quantity(),
request.unitPrice()
);
UUID orderId = commandBus.execute(command);
return ResponseEntity.status(HttpStatus.CREATED).body(orderId);
}
@GetMapping("/{id}")
public ResponseEntity
<OrderResponseDto> getOrder(@PathVariable UUID id) {
GetOrderQuery query = new GetOrderQuery(id);
OrderResponseDto response = queryBus.execute(query);
return ResponseEntity.ok(response);
}
}
5. Xử Lý Tranh Chấp Và Cập Nhật Chậm (Eventual Consistency Management)

⚠️ Verification Risk: Vấn đề nhất quán dữ liệu
Khi người dùng tạo đơn hàng và reload ngay lập tức, dữ liệu có thể chưa kịp đồng bộ từ PostgreSQL sang Redis. Đây là vấn đề kinh điển của Eventual Consistency.
Các giải pháp thực tế:
1. Client-side Retry với Exponential Backoff
async function getOrderWithRetry(orderId, maxRetries = 3) {
for (let i = 0; i < maxRetries; i++) {
try {
const response = await fetch(`/api/orders/${orderId}`);
if (response.ok) return await response.json();
if (response.status === 404 && i < maxRetries - 1) {
await sleep(Math.pow(2, i) * 100); // 100ms, 200ms, 400ms
continue;
}
throw new Error(`HTTP ${response.status}`);
} catch (e) {
if (i === maxRetries - 1) throw e;
await sleep(Math.pow(2, i) * 100);
}
}
}
2. Synchronous Event Processing (bỏ @Async)
Loại bỏ @Async để Event Listener xử lý đồng bộ trong cùng transaction. Điều này đảm bảo dữ liệu được đồng bộ trước khi response trả về client.
// Không dùng @Async - xử lý đồng bộ
@EventListener // Bỏ @Async
public void handleOrderCreated(OrderCreatedEvent event) {
readRepository.save(OrderView.fromEvent(event));
}
⚠️ Trade-off: Tăng thời gian response (do phải chờ Redis write) nhưng đảm bảo nhất quán dữ liệu ngay lập tức.
3. Idempotent Consumer + Dead Letter Queue (đã implement ở phần 4.10)
Trong code của chúng ta, chúng ta đã sử dụng setIfAbsent để đảm bảo idempotent và log lỗi để có thể gửi vào DLQ.
💡 Best Practice: Trong môi trường production, bạn nên tích hợp với Kafka hoặc RabbitMQ để có DLQ thực thụ và cơ chế retry với exponential backoff.
6. Lỗi Thường Gặp Khi Triển Khai CQRS
❌ Lỗi 1: Áp dụng CQRS cho ứng dụng CRUD đơn giản
Nguyên nhân: Tăng độ phức tạp mã nguồn gấp đôi mà không đem lại lợi ích rõ rệt.
Cách khắc phục: Chỉ áp dụng CQRS khi tỉ lệ Read/Write lệch nhau lớn (VD: 100:1) hoặc domain logic phức tạp.
❌ Lỗi 2: Đồng bộ dữ liệu không có cơ chế retry
Nguyên nhân: Khi Event Listener thất bại (Redis down, network timeout), dữ liệu sẽ không bao giờ được đồng bộ.
Cách khắc phục: Luôn implement retry mechanism hoặc Dead Letter Queue cho Event Listener (xem phần 4.10).
❌ Lỗi 3: Không xử lý idempotent cho Event
Nguyên nhân: Cùng một event có thể được xử lý nhiều lần (do retry, duplicate message).
Cách khắc phục: Sử dụng setIfAbsent với Redis hoặc unique constraint để đảm bảo idempotent.
❌ Lỗi 4: Query Model bị stale do không có cache invalidation
Nguyên nhân: Dữ liệu trong Redis không được cập nhật khi có thay đổi.
Cách khắc phục: Sử dụng TTL hợp lý và luôn update Redis khi có event update/delete.
❌ Lỗi 5: Quên bật @EnableAsync
Nguyên nhân: Sử dụng @Async nhưng không có @EnableAsync trong configuration.
Cách khắc phục: Luôn thêm @EnableAsync vào @SpringBootApplication hoặc config class.
7. Best Practices
✅ Sử dụng Java Records cho Command và Query DTOs để đảm bảo immutable và giảm boilerplate.
✅ Đảm bảo Event Listener có cơ chế Idempotent và Retry để tránh duplicate processing và thất bại không phục hồi.
✅ Sử dụng @TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT) thay vì @EventListener thông thường để đảm bảo event chỉ xử lý sau khi transaction commit.
✅ Đặt TTL hợp lý cho Redis Read Model để tránh dữ liệu stale tồn tại quá lâu.
✅ Tách biệt hoàn toàn Command và Query packages để dễ dàng bảo trì và mở rộng.
✅ Sử dụng Factory Methods trong Domain Entity để enforce business rules ngay từ khi tạo object.
✅ Logging chi tiết để dễ dàng debug khi có sự cố đồng bộ dữ liệu.
✅ Thêm Monitoring với Micrometer để theo dõi số lượng events, latency, và lỗi.
// Ví dụ monitoring với Micrometer
@Component
public class OrderEventListener {
private final MeterRegistry meterRegistry;
@Async
@TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT)
public void handleOrderCreated(OrderCreatedEvent event) {
meterRegistry.counter("cqrs.event.processed", "type", "order_created").increment();
long start = System.currentTimeMillis();
try {
// ... xử lý ...
meterRegistry.timer("cqrs.event.duration", "type", "order_created")
.record(System.currentTimeMillis() - start, TimeUnit.MILLISECONDS);
} catch (Exception e) {
meterRegistry.counter("cqrs.event.error", "type", "order_created").increment();
throw e;
}
}
}
8. FAQ
1. CQRS có bắt buộc phải đi kèm với Event Sourcing không?
Không. CQRS và Event Sourcing là hai pattern độc lập. Bạn có thể triển khai CQRS mà không cần Event Sourcing. Event Sourcing là một cách để implement CQRS nhưng không phải là bắt buộc. Trong bài viết này, chúng ta sử dụng Domain Events (qua ApplicationEventPublisher) mà không dùng Event Sourcing.
2. Làm sao xử lý việc người dùng reload trang ngay lập tức mà dữ liệu ở Read Model chưa kịp sync?
Có 3 cách tiếp cận chính:
- Client-side retry: Front-end tự động retry với exponential backoff
- Synchronous event processing: Bỏ
@Async, xử lý đồng bộ trong transaction - Optimistic UI: Hiển thị dữ liệu tạm (optimistic) và cập nhật sau khi có response
Lựa chọn nào phụ thuộc vào yêu cầu về trải nghiệm người dùng và hiệu năng của hệ thống.
3. Khi nào nên dùng message broker (Kafka/RabbitMQ) thay vì ApplicationEventPublisher?
- ApplicationEventPublisher: Phù hợp cho các ứng dụng đơn service, muốn giữ đơn giản, không cần phân tán.
- Kafka/RabbitMQ: Cần thiết khi có nhiều service, yêu cầu durable messaging, phân tán, hoặc cần replay event.
4. Có thể dùng CQRS với cùng một database không?
Có. CQRS chỉ yêu cầu tách biệt mô hình (model), không nhất thiết phải tách biệt database. Bạn có thể dùng cùng một database nhưng với các table/view/repository khác nhau cho read và write.
5. Làm thế nào để test CQRS?
- Unit Test: Test riêng Command Handler và Query Handler với Mockito
- Integration Test: Sử dụng
@SpringBootTestvà@TestTransactionalđể test luồng end-to-end - Contract Test: Sử dụng Spring Cloud Contract hoặc Pact để đảm bảo API response đúng format
Ví dụ Integration Test đơn giản:
@SpringBootTest
@AutoConfigureMockMvc
class OrderControllerTest {
@Autowired
private MockMvc mockMvc;
@Test
void createOrder_shouldReturnOrderId() throws Exception {
String request = """
{
"customerId": "cust123",
"productId": "prod456",
"quantity": 2,
"unitPrice": 49.99
}
""";
mockMvc.perform(post("/api/orders")
.contentType(MediaType.APPLICATION_JSON)
.content(request))
.andExpect(status().isCreated())
.andExpect(jsonPath("$").isNotEmpty());
}
}
9. Kết Luận
CQRS là một pattern kiến trúc mạnh mẽ nhưng không phải là giải pháp cho mọi bài toán. Trong bài viết này, chúng ta đã cùng nhau xây dựng một Lightweight CQRS hoàn chỉnh trong Spring Boot 3 sử dụng:
- Spring Data JPA cho Write Side (PostgreSQL)
- Spring Data Redis cho Read Side (Redis)
- Spring ApplicationEventPublisher để đồng bộ dữ liệu giữa hai model
Điểm mạnh của cách tiếp cận này là nhẹ, không phụ thuộc vào framework nặng như Axon, dễ hiểu và dễ bảo trì.
Hãy nhớ: Chỉ áp dụng CQRS khi bạn thực sự cần. Nếu tỷ lệ đọc/ghi của bạn không quá chênh lệch, một model CRUD truyền thống vẫn là lựa chọn tốt hơn.