Testcontainers로 Spring Boot Kafka 리스너 테스트하기

Testcontainers로 Spring Boot Kafka 리스너 테스트하기

이 가이드에서는 MySQL에 데이터를 저장하는 Kafka 리스너를 가진 Spring Boot 애플리케이션을 만들고, Testcontainers Kafka·MySQL 모듈과 Awaitility로 테스트하는 방법을 배워요.

출처: 문서

본문

이 가이드를 통해 다음 내용을 배울 수 있어요.

  • Kafka 통합이 포함된 Spring Boot 애플리케이션 만들기
  • Kafka 리스너를 구현하고 MySQL 데이터베이스에 데이터 저장하기
  • Testcontainers와 Awaitility로 Kafka 리스너 테스트하기

사전 준비 (Prerequisites)

  • Java 17 이상
  • Maven 또는 Gradle
  • Testcontainers가 지원하는 Docker 환경

참고: Testcontainers가 처음이라면 Testcontainers 개요를 방문해 알아보는 걸 권장해요.

Spring Boot 프로젝트 만들기

Spring Initializr에서 Spring for Apache Kafka, Spring Data JPA, MySQL Driver, Testcontainers 스타터를 선택해 Spring Boot 프로젝트를 만들어요. 또는 가이드 저장소를 클론해도 돼요. 애플리케이션을 생성한 뒤 Awaitility 라이브러리를 테스트 의존성으로 추가해요. 나중에 비동기 프로세스 흐름의 기대값을 단언하는 데 사용할 거예요. pom.xml의 핵심 의존성은 다음과 같아요.

<properties>
    <java.version>17</java.version>
    <testcontainers.version>2.0.4</testcontainers.version>
</properties>

<dependencies>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-data-jpa</artifactId>
    </dependency>
    <dependency>
        <groupId>org.springframework.kafka</groupId>
        <artifactId>spring-kafka</artifactId>
    </dependency>
    <dependency>
        <groupId>com.mysql</groupId>
        <artifactId>mysql-connector-j</artifactId>
        <scope>runtime</scope>
    </dependency>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-test</artifactId>
        <scope>test</scope>
    </dependency>
    <dependency>
        <groupId>org.springframework.kafka</groupId>
        <artifactId>spring-kafka-test</artifactId>
        <scope>test</scope>
    </dependency>
    <dependency>
        <groupId>org.testcontainers</groupId>
        <artifactId>testcontainers-junit-jupiter</artifactId>
        <scope>test</scope>
    </dependency>
    <dependency>
        <groupId>org.testcontainers</groupId>
        <artifactId>testcontainers-kafka</artifactId>
        <scope>test</scope>
    </dependency>
    <dependency>
        <groupId>org.testcontainers</groupId>
        <artifactId>testcontainers-mysql</artifactId>
        <scope>test</scope>
    </dependency>
    <dependency>
        <groupId>org.awaitility</groupId>
        <artifactId>awaitility</artifactId>
        <scope>test</scope>
    </dependency>
</dependencies>

모든 Testcontainers 모듈 의존성에 버전을 반복하지 않도록 Testcontainers BOM(Bill of Materials)을 사용하는 걸 권장해요.

JPA 엔티티 만들기

애플리케이션은 product-price-changes라는 토픽을 구독해요. 메시지가 도착하면 이벤트 페이로드에서 제품 코드와 가격을 추출해 MySQL 데이터베이스에서 해당 제품의 가격을 업데이트해요. Product.java를 만들어요.

package com.testcontainers.demo;

import jakarta.persistence.Column;
import jakarta.persistence.Entity;
import jakarta.persistence.GeneratedValue;
import jakarta.persistence.GenerationType;
import jakarta.persistence.Id;
import jakarta.persistence.Table;
import java.math.BigDecimal;

@Entity
@Table(name = "products")
class Product {
    @Id
    @GeneratedValue(strategy = GenerationType.IDENTITY)
    private Long id;

    @Column(nullable = false, unique = true)
    private String code;

    @Column(nullable = false)
    private String name;

    @Column(nullable = false)
    private BigDecimal price;

    public Product() {}

    public Product(Long id, String code, String name, BigDecimal price) {
        this.id = id;
        this.code = code;
        this.name = name;
        this.price = price;
    }

    public Long getId() { return id; }
    public void setId(Long id) { this.id = id; }

    public String getCode() { return code; }
    public void setCode(String code) { this.code = code; }

    public String getName() { return name; }
    public void setName(String name) { this.name = name; }

    public BigDecimal getPrice() { return price; }
    public void setPrice(BigDecimal price) { this.price = price; }
}

Spring Data JPA 리포지토리 만들기

Product 엔티티에 대한 리포지토리 인터페이스를 만들고, 코드로 제품을 찾는 메서드와 주어진 제품 코드의 가격을 업데이트하는 메서드를 추가해요.

package com.testcontainers.demo;

import java.math.BigDecimal;
import java.util.Optional;

import org.springframework.data.jpa.repository.JpaRepository;
import org.springframework.data.jpa.repository.Modifying;
import org.springframework.data.jpa.repository.Query;
import org.springframework.data.repository.query.Param;

interface ProductRepository extends JpaRepository<Product, Long> {
    Optional<Product> findByCode(String code);

    @Modifying
    @Query("update Product p set p.price = :price where p.code = :productCode")
    void updateProductPrice(@Param("productCode") String productCode, @Param("price") BigDecimal price);
}

스키마 생성 스크립트 추가하기

애플리케이션이 인메모리 데이터베이스를 사용하지 않으므로 MySQL 테이블을 직접 만들어야 해요. 프로덕션에서 권장하는 방법은 Flyway나 Liquibase 같은 마이그레이션 도구지만, 이 가이드에서는 Spring Boot 내장 스키마 초기화로 충분해요. src/main/resources/schema.sql을 만들어요.

create table products (
    id int NOT NULL AUTO_INCREMENT,
    code varchar(255) not null,
    name varchar(255) not null,
    price numeric(5, 2) not null,
    PRIMARY KEY (id),
    UNIQUE (code)
);

src/main/resources/application.properties에서 스키마 초기화를 활성화해요.

spring.sql.init.mode=always

이벤트 페이로드 만들기

Kafka 토픽에서 받는 이벤트 페이로드의 구조를 나타내는 ProductPriceChangedEvent 레코드를 만들어요.

package com.testcontainers.demo;

import java.math.BigDecimal;

record ProductPriceChangedEvent(String productCode, BigDecimal price) {}

발신자와 수신자는 다음 JSON 형식에 동의해요.

{
  "productCode": "P100",
  "price": 25.0
}

Kafka 리스너 구현하기

product-price-changes 토픽의 메시지를 처리하고 데이터베이스에서 제품 가격을 업데이트하는 ProductPriceChangedEventHandler.java를 만들어요.

package com.testcontainers.demo;

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component;
import org.springframework.transaction.annotation.Transactional;

@Component
@Transactional
class ProductPriceChangedEventHandler {
    private static final Logger log = LoggerFactory.getLogger(ProductPriceChangedEventHandler.class);

    private final ProductRepository productRepository;

    ProductPriceChangedEventHandler(ProductRepository productRepository) {
        this.productRepository = productRepository;
    }

    @KafkaListener(topics = "product-price-changes", groupId = "demo")
    public void handle(ProductPriceChangedEvent event) {
        log.info("Received a ProductPriceChangedEvent with productCode:{}: ", event.productCode());
        productRepository.updateProductPrice(event.productCode(), event.price());
    }
}

@KafkaListener 어노테이션은 들을 토픽 이름을 지정해요. Spring Kafka는 application.properties에 구성된 속성에 따라 직렬화와 역직렬화를 처리해요.

Kafka 직렬화 구성하기

src/main/resources/application.properties에 다음 Kafka 속성을 추가해요.

######## Kafka Configuration #########
spring.kafka.bootstrap-servers=localhost:9092
spring.kafka.producer.key-serializer=org.apache.kafka.common.serialization.StringSerializer
spring.kafka.producer.value-serializer=org.springframework.kafka.support.serializer.JsonSerializer
spring.kafka.consumer.group-id=demo
spring.kafka.consumer.auto-offset-reset=latest
spring.kafka.consumer.key-deserializer=org.apache.kafka.common.serialization.StringDeserializer
spring.kafka.consumer.value-deserializer=org.springframework.kafka.support.serializer.JsonDeserializer
spring.kafka.consumer.properties.spring.json.trusted.packages=com.testcontainers.demo

productCode 키는 StringSerializer/StringDeserializer로 (역)직렬화되고, ProductPriceChangedEvent 값은 JsonSerializer/JsonDeserializer로 (역)직렬화돼요.

Testcontainers로 테스트 작성하기

Kafka 리스너를 테스트하려면 실행 중인 Kafka 브로커와 MySQL 데이터베이스, 그리고 시작된 Spring 컨텍스트가 필요해요. Testcontainers가 두 서비스를 Docker 컨테이너에서 띄우고 @DynamicPropertySource가 이를 Spring에 연결해요.

ProductPriceChangedEventHandlerTest.java를 만들어요.

package com.testcontainers.demo;

import static java.util.concurrent.TimeUnit.SECONDS;
import static org.assertj.core.api.Assertions.assertThat;
import static org.awaitility.Awaitility.await;

import java.math.BigDecimal;
import java.time.Duration;
import java.util.Optional;

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.kafka.core.KafkaTemplate;
import org.springframework.test.context.DynamicPropertyRegistry;
import org.springframework.test.context.DynamicPropertySource;
import org.springframework.test.context.TestPropertySource;
import org.testcontainers.kafka.ConfluentKafkaContainer;
import org.testcontainers.junit.jupiter.Container;
import org.testcontainers.junit.jupiter.Testcontainers;

@SpringBootTest
@TestPropertySource(properties = {
    "spring.kafka.consumer.auto-offset-reset=earliest",
    "spring.datasource.url=jdbc:tc:mysql:8.0.32:///db",
})
@Testcontainers
class ProductPriceChangedEventHandlerTest {
    @Container
    static final ConfluentKafkaContainer kafka = new ConfluentKafkaContainer("confluentinc/cp-kafka:7.8.0");

    @DynamicPropertySource
    static void overrideProperties(DynamicPropertyRegistry registry) {
        registry.add("spring.kafka.bootstrap-servers", kafka::getBootstrapServers);
    }

    @Autowired
    private KafkaTemplate<String, Object> kafkaTemplate;

    @Autowired
    private ProductRepository productRepository;

    @BeforeEach
    void setUp() {
        Product product = new Product(null, "P100", "Product One", BigDecimal.TEN);
        productRepository.save(product);
    }

    @Test
    void shouldHandleProductPriceChangedEvent() {
        ProductPriceChangedEvent event = new ProductPriceChangedEvent("P100", new BigDecimal("14.50"));
        kafkaTemplate.send("product-price-changes", event.productCode(), event);

        await()
            .pollInterval(Duration.ofSeconds(3))
            .atMost(10, SECONDS)
            .untilAsserted(() -> {
                Optional<Product> optionalProduct = productRepository.findByCode("P100");
                assertThat(optionalProduct).isPresent();
                assertThat(optionalProduct.get().getCode()).isEqualTo("P100");
                assertThat(optionalProduct.get().getPrice())
                    .isEqualTo(new BigDecimal("14.50"));
            });
    }
}

테스트가 하는 일을 살펴보면,

  • @SpringBootTest는 전체 Spring 애플리케이션 컨텍스트를 시작해요.
  • @TestPropertySource의 Testcontainers 특수 JDBC URL(jdbc:tc:mysql:8.0.32:///db)이 MySQL 컨테이너를 띄우고 데이터소스로 자동 구성해요.
  • @Testcontainers와 @Container가 Kafka 컨테이너 수명주기를 관리해요.
  • @DynamicPropertySource는 Kafka 부트스트랩 서버를 Spring에 등록해서 프로듀서와 컨슈머가 테스트 컨테이너에 연결되게 해요.
  • @BeforeEach는 각 테스트 전에 데이터베이스에 Product 레코드를 만들어요.
  • 테스트는 KafkaTemplate으로 product-price-changes 토픽에 ProductPriceChangedEvent를 전송해요. Spring Boot가 JsonSerializer로 객체를 JSON으로 변환해요.
  • Kafka 메시지 처리는 비동기이므로 테스트는 Awaitility를 사용해 데이터베이스의 제품 가격이 기대값과 일치할 때까지 3초마다(최대 10초) 폴링해요.
  • spring.kafka.consumer.auto-offset-reset 속성을 earliest로 설정하면, 리스너가 준비되기 전에 토픽에 보내진 메시지라도 소비하게 해요. 이 설정은 테스트 실행 시 유용해요.

테스트 실행과 다음 단계

테스트를 실행해요.

$ ./mvnw test

또는 Gradle로,

$ ./gradlew test

Kafka와 MySQL Docker 컨테이너가 시작되고 모든 테스트가 통과하는 걸 볼 수 있어요. 테스트가 끝나면 컨테이너는 자동으로 중지되고 제거돼요.

요약 (Summary)

목이나 내장 대안보다 실제 Kafka와 MySQL 인스턴스로 테스트하면 코드의 정확성에 더 큰 확신을 얻을 수 있어요. Testcontainers 라이브러리가 컨테이너 수명주기를 관리하므로, 통합 테스트가 프로덕션에서 쓰는 것과 같은 서비스에 대해 실행돼요.

Testcontainers에 대해 더 알아보고 싶다면 Testcontainers 개요를 방문해요.

더 읽어보기 (Further reading)

  • Java Spring Boot 프로젝트에서 Testcontainers 시작하기
  • 테스트용 H2를 실제 데이터베이스로 교체하는 가장 간단한 방법
  • Awaitility
  • Testcontainers Kafka 모듈
  • Testcontainers MySQL 모듈

더 알아보기 (Learn more)