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 모듈