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