Testcontainers로 Micronaut Kafka 리스너 테스트하기

Testcontainers로 Micronaut Kafka 리스너 테스트하기

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

출처: 문서

본문

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

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

사전 준비 (Prerequisites)

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

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

Micronaut 프로젝트 만들기

Micronaut Launch에서 kafka, data-jpa, mysql, awaitility, assertj, testcontainers 기능을 선택해 Micronaut 프로젝트를 만들어요. 또는 가이드 저장소를 클론해도 돼요.

비동기 프로세스 흐름의 기대값을 단언하기 위해 Awaitility 라이브러리를 사용해요. pom.xml의 핵심 의존성은 다음과 같아요.

<parent>
    <groupId>io.micronaut.platform</groupId>
    <artifactId>micronaut-parent</artifactId>
    <version>4.1.4</version>
</parent>

<dependencies>
    <dependency>
        <groupId>io.micronaut.data</groupId>
        <artifactId>micronaut-data-hibernate-jpa</artifactId>
        <scope>compile</scope>
    </dependency>
    <dependency>
        <groupId>io.micronaut.kafka</groupId>
        <artifactId>micronaut-kafka</artifactId>
        <scope>compile</scope>
    </dependency>
    <dependency>
        <groupId>io.micronaut.serde</groupId>
        <artifactId>micronaut-serde-jackson</artifactId>
        <scope>compile</scope>
    </dependency>
    <dependency>
        <groupId>io.micronaut.sql</groupId>
        <artifactId>micronaut-jdbc-hikari</artifactId>
        <scope>compile</scope>
    </dependency>
    <dependency>
        <groupId>mysql</groupId>
        <artifactId>mysql-connector-java</artifactId>
        <scope>runtime</scope>
    </dependency>
    <dependency>
        <groupId>org.awaitility</groupId>
        <artifactId>awaitility</artifactId>
        <version>4.2.0</version>
        <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>
</dependencies>

Micronaut 부모 POM이 Testcontainers BOM을 관리하므로, Testcontainers 모듈별로 버전을 따로 지정할 필요가 없어요.

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")
public 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; }
}

Micronaut Data JPA 리포지토리 만들기

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

package com.testcontainers.demo;

import io.micronaut.data.annotation.Query;
import io.micronaut.data.annotation.Repository;
import io.micronaut.data.jpa.repository.JpaRepository;
import java.math.BigDecimal;
import java.util.Optional;

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

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

Spring Data JPA와 달리 Micronaut Data는 컴파일 타임 어노테이션 처리를 사용해 리포지토리 메서드를 구현하므로 런타임 리플렉션을 피할 수 있어요.

이벤트 페이로드 만들기

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

package com.testcontainers.demo;

import io.micronaut.serde.annotation.Serdeable;
import java.math.BigDecimal;

@Serdeable
public record ProductPriceChangedEvent(String productCode, BigDecimal price) {}

@Serdeable 어노테이션은 Micronaut Serialization에게 이 타입이 직렬화·역직렬화될 수 있음을 알려줘요. 발신자와 수신자는 다음 JSON 형식에 동의해요.

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

Kafka 리스너 구현하기

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

package com.testcontainers.demo;

import static io.micronaut.configuration.kafka.annotation.OffsetReset.EARLIEST;

import io.micronaut.configuration.kafka.annotation.KafkaListener;
import io.micronaut.configuration.kafka.annotation.Topic;
import jakarta.inject.Singleton;
import jakarta.transaction.Transactional;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

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

    private final ProductRepository productRepository;

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

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

주요 세부 사항:

  • @KafkaListener 어노테이션이 이 클래스를 Kafka 메시지 리스너로 지정해요.
  • offsetReset을 EARLIEST로 설정하면 리스너가 파티션의 처음부터 메시지 소비를 시작하는데, 테스트 중에 유용해요.
  • @Topic 어노테이션은 구독할 토픽을 지정해요.
  • Micronaut는 Micronaut Serialization을 사용해 ProductPriceChangedEvent의 JSON 역직렬화를 자동으로 처리해요.

데이터소스 구성하기

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

micronaut.application.name=tc-guide-testing-micronaut-kafka-listener
datasources.default.db-type=mysql
datasources.default.dialect=MYSQL
jpa.default.properties.hibernate.hbm2ddl.auto=update
jpa.default.entity-scan.packages=com.testcontainers.demo
datasources.default.driver-class-name=com.mysql.cj.jdbc.Driver

Hibernate의 hbm2ddl.auto=update는 데이터베이스 스키마를 자동으로 생성하고 업데이트해요. 테스트에서는 테스트 속성 파일에서 이를 create-drop으로 덮어쓸 거예요. src/test/resources/application-test.properties를 만들어요.

jpa.default.properties.hibernate.hbm2ddl.auto=create-drop

Testcontainers로 테스트 작성하기

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

테스트용 Kafka 클라이언트 만들기

먼저 테스트에서 이벤트를 게시할 @KafkaClient 인터페이스를 만들어요.

package com.testcontainers.demo;

import io.micronaut.configuration.kafka.annotation.KafkaClient;
import io.micronaut.configuration.kafka.annotation.KafkaKey;
import io.micronaut.configuration.kafka.annotation.Topic;

@KafkaClient
public interface ProductPriceChangesClient {
    @Topic("product-price-changes")
    void send(@KafkaKey String productCode, ProductPriceChangedEvent event);
}

주요 세부 사항:

  • @KafkaClient 어노테이션은 이 인터페이스를 Kafka 프로듀서로 지정해요.
  • @Topic 어노테이션은 대상 토픽을 지정해요.
  • @KafkaKey 어노테이션은 Kafka 메시지 키로 사용될 매개변수를 표시해요. 그런 매개변수가 없으면 Micronaut는 null 키로 레코드를 보내요.

테스트 작성하기

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 io.micronaut.context.annotation.Property;
import io.micronaut.core.annotation.NonNull;
import io.micronaut.test.extensions.junit5.annotation.MicronautTest;
import io.micronaut.test.support.TestPropertyProvider;
import java.math.BigDecimal;
import java.time.Duration;
import java.util.Collections;
import java.util.Map;
import java.util.Optional;

import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.TestInstance;
import org.testcontainers.kafka.ConfluentKafkaContainer;
import org.testcontainers.junit.jupiter.Container;
import org.testcontainers.junit.jupiter.Testcontainers;

@MicronautTest(transactional = false)
@Property(name = "datasources.default.driver-class-name", value = "org.testcontainers.jdbc.ContainerDatabaseDriver")
@Property(name = "datasources.default.url", value = "jdbc:tc:mysql:8.0.32:///db")
@Testcontainers(disabledWithoutDocker = true)
@TestInstance(TestInstance.Lifecycle.PER_CLASS)
class ProductPriceChangedEventHandlerTest implements TestPropertyProvider {
    @Container
    static final ConfluentKafkaContainer kafka = new ConfluentKafkaContainer("confluentinc/cp-kafka:7.8.0");

    @Override
    public @NonNull Map<String, String> getProperties() {
        if (!kafka.isRunning()) {
            kafka.start();
        }
        return Collections.singletonMap("kafka.bootstrap.servers", kafka.getBootstrapServers());
    }

    @Test
    void shouldHandleProductPriceChangedEvent(
            ProductPriceChangesClient productPriceChangesClient,
            ProductRepository productRepository) {
        Product product = new Product(null, "P100", "Product One", BigDecimal.TEN);
        Long id = productRepository.save(product).getId();

        ProductPriceChangedEvent event = new ProductPriceChangedEvent("P100", new BigDecimal("14.50"));
        productPriceChangesClient.send(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"));
        });

        productRepository.deleteById(id);
    }
}

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

  • @MicronautTest는 Micronaut 애플리케이션 컨텍스트와 내장 서버를 초기화해요.
  • transactional을 false로 설정하면 각 테스트 메서드가 롤백되는 트랜잭션 안에서 실행되지 않도록 해요. Kafka 리스너는 별도 스레드에서 메시지를 처리하기 때문에 반드시 필요해요.
  • @Property 어노테이션은 데이터소스 드라이버와 URL을 Testcontainers 특수 JDBC URL(jdbc:tc:mysql:8.0.32:///db)로 덮어써요. 이 URL이 MySQL 컨테이너를 띄우고 자동으로 데이터소스로 구성해요.
  • @Testcontainers와 @Container가 Kafka 컨테이너 수명주기를 관리해요.
  • TestPropertyProvider 인터페이스는 Kafka 부트스트랩 서버를 Micronaut에 등록해서, 프로듀서와 컨슈머가 테스트 컨테이너에 연결되게 해요.
  • @TestInstance(TestInstance.Lifecycle.PER_CLASS)는 모든 테스트 메서드에 대해 단일 테스트 인스턴스를 만들어요. TestPropertyProvider를 구현할 때 필요해요.
  • 테스트는 데이터베이스에 Product 레코드를 만든 다음, ProductPriceChangesClient로 product-price-changes 토픽에 ProductPriceChangedEvent를 전송해요.
  • Kafka 메시지 처리는 비동기이므로 테스트는 Awaitility를 사용해 데이터베이스의 제품 가격이 기대값과 일치할 때까지 3초마다(최대 10초) 폴링해요.

테스트 실행과 다음 단계

테스트를 실행해요.

$ ./mvnw test

또는 Gradle로,

$ ./gradlew test

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

요약 (Summary)

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

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

더 읽어보기 (Further reading)

  • Micronaut 앱에서 WireMock으로 REST API 통합 테스트하기
  • Testcontainers로 Spring Boot Kafka 리스너 테스트하기
  • Java Spring Boot 프로젝트에서 Testcontainers 시작하기
  • Awaitility
  • Testcontainers Kafka 모듈
  • Testcontainers MySQL 모듈

더 알아보기 (Learn more)