포스트

Kafka를 활용한 CQRS (3) - Spring Boot로 CQRS 구현하기

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 처리 (데이터 조회)
저장소MariaDBMongoDB
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

데이터베이스와 연동하는 BookRepositoryJpaRepository를 상속받는 인터페이스로 간단히 구성한다.

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.gradledependencies에 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에서 ProducerFactoryKafkaTemplate을 빈으로 등록한다.

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());
    }
}

읽기 서비스에서 저장된 책이 무엇인지 식별할 수 있어야 하므로, BookDTObid 필드를 추가한다. 그리고 메시지를 실제로 전송하는 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.ymlKafkaConfiguration은 쓰기 서비스와 동일하게 맞춘다. 그 위에 @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로 이어봤다. 쓰기 서비스는 책을 저장한 뒤 KafkaProducerbid를 담은 메시지를 발행하고, 읽기 서비스는 @KafkaListener로 이 메시지를 구독해서 자신의 MongoDB에 반영한다. 이렇게 두 서비스가 데이터베이스를 직접 공유하지 않고 Kafka 메시지만으로 느슨하게 연결되는 구조가, 앞서 정리한 CQRS의 “쓰기와 읽기의 책임 분리”를 코드 수준에서 구현한 형태라고 볼 수 있다.

이 기사는 저작권자의 CC BY 4.0 라이센스를 따릅니다.