Kafka를 활용한 CQRS (3) - Spring Boot로 CQRS 구현하기
1. 개요
앞선 글에서 CQRS의 개념과 Kafka의 기초를 정리했으니, 이번 글에서는 Spring Boot로 실제 쓰기(Command)/읽기(Query) 서비스를 나눠 만들고 Kafka로 두 서비스를 연결하는 예제를 구현한다. 쓰기 서비스(cqrs_write)는 MariaDB에 도서 데이터를 저장하는 역할을, 읽기 서비스(cqrs_read)는 MongoDB에서 도서 데이터를 조회하는 역할을 각각 전담하며, 쓰기 작업이 끝났다는 사실을 Kafka 메시지로 알려주면 읽기 서비스가 이를 구독해서 자신의 저장소(MongoDB)에 반영하는 흐름을 만드는 것이 목표다.
2. 전체 구조
두 서비스는 이렇게 분리된다.
1
2
3
4
5
6
7
8
9
클라이언트 → (POST) → 쓰기 서비스(cqrs_write, MariaDB)
│
└─ (Kafka에 "책 저장됨" 메시지 발행: bid)
│
▼
읽기 서비스(cqrs_read, MongoDB)
(Kafka 메시지 구독 후 MongoDB에 반영)
▲
클라이언트 ← (GET) ←─────────────────┘
| 구분 | 쓰기 서비스 (cqrs_write) | 읽기 서비스 (cqrs_read) |
|---|---|---|
| 역할 | Command 처리 (데이터 저장) | Query 처리 (데이터 조회) |
| 저장소 | MariaDB | MongoDB |
| Kafka 역할 | 저장 완료 후 메시지 프로듀서 | 메시지를 받아 저장소에 반영하는 컨슈머 |
3. 쓰기 서비스 구현 (cqrs_write)
3.1 프로젝트 생성과 설정
Spring Dev Tools, Lombok, Spring Web, Spring Data JPA, MariaDB, Kafka 의존성을 가지고 프로젝트를 생성한다. application.yml에는 서버 포트와 MariaDB 접속 정보를 지정한다.
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
server:
port: 8080
spring:
config:
import: optional:file:.env[.properties]
application:
name: autoeverspring
datasource:
url: ${db_url}
driver-class-name: org.mariadb.jdbc.Driver
username: ${username}
password: ${password}
jpa:
hibernate:
ddl-auto: update
properties:
hibernate:
format_sql: true
show_sql: true
kafka:
bootstrap-servers: 127.0.0.1:9092
consumer:
group-id: adamsoft
auto-offset-reset: earliest
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
producer:
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.apache.kafka.common.serialization.StringSerializer
logging:
level:
org.hibernate.type.description.sql: trace
3.2 Entity와 DTO
Book은 데이터베이스 테이블과 직접 연결되는 Entity 클래스다.
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
import jakarta.persistence.*;
import lombok.*;
import java.util.Date;
@Entity
@Table(name = "book")
@ToString
@Getter
@Builder
@AllArgsConstructor
@NoArgsConstructor
public class Book {
@Id
@GeneratedValue(strategy = GenerationType.AUTO)
private Long bid;
@Column(length = 50, nullable = false)
private String title;
@Column(length = 50, nullable = false)
private String author;
@Column(length = 50, nullable = false)
private String category;
@Column
private int pages;
@Column
private int price;
@Column
private Date published_date;
@Column(length = 50, nullable = false)
private String description;
}
BookDTO는 계층 간 데이터 이동을 위한 클래스다.
1
2
3
4
5
6
7
8
9
10
11
12
import lombok.Data;
@Data
public class BookDTO {
private String title;
private String author;
private String category;
private int pages;
private int price;
private String published_date;
private String description;
}
3.3 Repository, Service, Controller
데이터베이스와 연동하는 BookRepository는 JpaRepository를 상속받는 인터페이스로 간단히 구성한다.
1
2
3
4
import org.springframework.data.jpa.repository.JpaRepository;
public interface BookRepository extends JpaRepository<Book, Long> {
}
BookService는 DTO를 Entity로 변환해서 저장하는 역할을 한다.
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
import lombok.RequiredArgsConstructor;
import org.springframework.stereotype.Service;
import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.Locale;
@Service
@RequiredArgsConstructor
public class BookService {
private final BookRepository bookRepository;
public void saveBook(BookDTO bookDTO){
try {
SimpleDateFormat formatter = new SimpleDateFormat(
"yyyy-MM-dd", Locale.ENGLISH);
Date published_date = formatter.parse(bookDTO.getPublished_date());
Book book = Book.builder()
.title(bookDTO.getTitle())
.author(bookDTO.getAuthor())
.category(bookDTO.getCategory())
.pages(bookDTO.getPages())
.price(bookDTO.getPrice())
.published_date(published_date)
.description(bookDTO.getDescription())
.build();
bookRepository.save(book);
}
catch(Exception e){
System.out.println(e.getMessage());
}
}
}
BookController는 사용자의 요청을 받아 서비스로 넘긴다.
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
import lombok.RequiredArgsConstructor;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RestController;
@RestController
@RequiredArgsConstructor
public class BookController {
private final BookService bookService;
@GetMapping("/")
public String index(){
return "homepage";
}
@GetMapping("/health")
public String healthCheck(){
return "success";
}
@PostMapping("/cqrs/book")
public String saveBook(@RequestBody BookDTO bookDTO){
bookService.saveBook(bookDTO);
return "success";
}
}
localhost:7000으로 웹 서버를 구동한 뒤, REST API 테스트에 유용한 POSTMAN으로 데이터 삽입을 테스트한다.
4. 읽기 서비스 구현 (cqrs_read)
Spring Boot Devtools, Lombok, Spring Web, Spring Data JPA, Spring Data MongoDB, Kafka 의존성으로 별도 프로젝트를 생성한다. 아직 Kafka를 연결하기 전이므로, BookController는 MongoDB에 직접 접속해서 저장된 도서 목록을 조회하는 역할만 한다.
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
import com.mongodb.client.*;
import lombok.RequiredArgsConstructor;
import org.bson.Document;
import org.springframework.http.HttpStatus;
import org.springframework.http.ResponseEntity;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RestController;
import java.util.ArrayList;
import java.util.List;
@RestController
@RequiredArgsConstructor
public class BookController {
@GetMapping("/cqrs/book")
public ResponseEntity<List> getBooks(){
MongoClient mongoClient = MongoClients.create("mongodb://localhost:27017");
MongoDatabase database = mongoClient.getDatabase("mymongo");
MongoCollection<Document> mongo_books =
database.getCollection("books");
List<Document> list = new ArrayList<Document>();
try{
try(MongoCursor<Document> cur = mongo_books.find().iterator()){
while(cur.hasNext()){
Document doc = cur.next();
list.add(doc);
}
}
}catch(Exception e){
e.printStackTrace();
}
finally{
mongoClient.close();
}
return ResponseEntity.status(HttpStatus.OK).body(list);
}
}
실행한 뒤 MongoDB에 미리 넣어둔 데이터가 정상적으로 조회되는지 확인한다.
5. Kafka로 두 서비스 연결하기
지금까지는 쓰기 서비스와 읽기 서비스가 각자의 저장소만 바라볼 뿐 서로 연결되어 있지 않다. 이 둘을 Kafka로 이어서, 쓰기 서비스에 책이 저장되면 그 사실이 읽기 서비스의 MongoDB에도 반영되도록 만든다.
5.1 쓰기 서비스에 Kafka Producer 추가
build.gradle의 dependencies에 Kafka 라이브러리를 추가하고 리로드한다.
1
implementation 'org.springframework.boot:spring-boot-starter-kafka'
application.yml에 Kafka 접속 정보를 추가한다.
1
2
3
4
5
6
7
8
9
10
11
spring:
kafka:
bootstrap-servers: 127.0.0.1:9092
consumer:
group-id: adamsoft
auto-offset-reset: earliest
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
producer:
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.apache.kafka.common.serialization.StringSerializer
KafkaConfiguration에서 ProducerFactory와 KafkaTemplate을 빈으로 등록한다.
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
@Configuration
public class KafkaConfiguration {
@Value("${spring.kafka.bootstrap-servers}")
private String bootStrapServers;
@Bean
public ProducerFactory<String, String> producerFactory(){
Map<String, Object> configs = new HashMap<>();
configs.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootStrapServers);
configs.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
configs.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
return new DefaultKafkaProducerFactory<>(configs);
}
@Bean
public KafkaTemplate<String, String> kafkaTemplate(){
return new KafkaTemplate<>(producerFactory());
}
}
읽기 서비스에서 저장된 책이 무엇인지 식별할 수 있어야 하므로, BookDTO에 bid 필드를 추가한다. 그리고 메시지를 실제로 전송하는 KafkaProducer를 작성한다.
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
@Service
@RequiredArgsConstructor
public class KafkaProducer {
private static final String TOPIC = "cqrs-topic";
@Autowired
private KafkaTemplate<String, String> kafkaTemplate;
private final Logger log = LoggerFactory.getLogger(getClass());
public void sendMessage(BookDTO bookDTO){
String message = "{\"bid\":" + "\"" + bookDTO.getBid() + "\"}";
//메시지 전송
this.kafkaTemplate.send(TOPIC, message);
}
}
마지막으로 BookService가 책을 저장한 직후 KafkaProducer를 호출해서 저장된 bid를 메시지로 전송하도록 수정한다.
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
@Service
@RequiredArgsConstructor
public class BookService {
private final BookRepository bookRepository;
private final KafkaProducer kafkaProducer;
public void saveBook(BookDTO bookDTO){
try {
SimpleDateFormat formatter = new SimpleDateFormat(
"yyyy-MM-dd", Locale.ENGLISH);
Date published_date = formatter.parse(bookDTO.getPublished_date());
Book book = Book.builder()
.title(bookDTO.getTitle())
.author(bookDTO.getAuthor())
.category(bookDTO.getCategory())
.pages(bookDTO.getPages())
.price(bookDTO.getPrice())
.published_date(published_date)
.description(bookDTO.getDescription())
.build();
bookRepository.save(book);
bookDTO.setBid(book.getBid());
//쓰기 작업을 완료할 때 카프카에게 메시지를 전송
kafkaProducer.sendMessage(bookDTO);
}
catch(Exception e){
System.out.println(e.getMessage());
}
}
}
5.2 토픽 생성과 메시지 확인
터미널에서 Kafka 컨테이너에 접속해 cqrs-topic을 만들고, 콘솔 컨슈머로 메시지가 잘 들어오는지 먼저 확인한다.
1
2
3
4
5
docker exec -it kafka /bin/bash
kafka-topics.sh --create --zookeeper zookeeper:2181 --replication-factor 1 --partitions 1 --topic cqrs-topic
kafka-topics.sh --bootstrap-server localhost:9092 --list
kafka-console-consumer.sh --topic cqrs-topic --bootstrap-server localhost:9092 --from-beginning
이 상태에서 POSTMAN으로 POST http://localhost:7000/cqrs/book에 아래와 같은 JSON을 담아 요청을 보내면, MariaDB에 데이터가 저장됨과 동시에 콘솔 컨슈머에 bid 메시지가 출력되는 것을 확인할 수 있다.
1
2
3
4
5
6
7
8
9
{
"title":"오딧세이",
"author":"호메로스",
"category":"신화",
"description":"크리스토퍼 놀란",
"pages":356,
"price":15000,
"published_date":"2026-08-13"
}
5.3 읽기 서비스에 Kafka Consumer 추가
읽기 서비스(cqrs_read)의 build.gradle에는 메시지를 JSON으로 파싱하기 위한 라이브러리를 추가한다.
1
implementation 'org.json:json:20190722'
application.yml과 KafkaConfiguration은 쓰기 서비스와 동일하게 맞춘다. 그 위에 @KafkaListener로 토픽을 구독하는 kafkaConsumer를 작성한다.
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
@Service
public class kafkaConsumer {
@KafkaListener(topics="cqrs-topic", groupId = "adamsoft")
public void consumer(String message) throws IOException{
System.out.println("message:" + message);
//JSON 파싱: JSON문자열로 온 것을 자바 객체로 변환
JSONObject messageObj = new JSONObject(message);
//MongoDB 컬렉션에 연결
MongoClient mongoClient = MongoClients.create("mongodb://localhost:27017");
MongoDatabase database = mongoClient.getDatabase("mymongo");
MongoCollection<Document> mongo_books =
database.getCollection("books");
//받은 데이터로 삽입할 데이터를 생성
Document book = new Document();
book.append("bid", messageObj.getLong("bid"));
mongo_books.insertOne(book);
mongoClient.close();
}
}
두 서비스를 함께 실행한 상태에서 쓰기 서비스에 책을 저장하면, Kafka를 거쳐 읽기 서비스가 자동으로 MongoDB에 같은 bid를 기록하는 것을 확인할 수 있다.
6. 정리
이번 글에서는 CQRS의 Command와 Query를 실제로 두 개의 Spring Boot 서비스(쓰기 서비스는 MariaDB, 읽기 서비스는 MongoDB)로 분리해서 구현하고, 그 사이를 Kafka로 이어봤다. 쓰기 서비스는 책을 저장한 뒤 KafkaProducer로 bid를 담은 메시지를 발행하고, 읽기 서비스는 @KafkaListener로 이 메시지를 구독해서 자신의 MongoDB에 반영한다. 이렇게 두 서비스가 데이터베이스를 직접 공유하지 않고 Kafka 메시지만으로 느슨하게 연결되는 구조가, 앞서 정리한 CQRS의 “쓰기와 읽기의 책임 분리”를 코드 수준에서 구현한 형태라고 볼 수 있다.