이 도메인의 학습 목표

이 도메인이 묻는 것

출제 형태는 사실상 두 가지입니다.

  1. 도구 선택 — "Streams 토폴로지 로직만 빠르게 검증하려면?", "실제 브로커 동작까지 확인하려면?" → matching 유형(도구 ↔ 용도)으로도 자주 나옵니다.
  2. 도구의 성질 — 브로커가 필요한지, 시간을 제어할 수 있는지, 무엇을 검증할 수 없는지.

참고 챕터는 10장의 테스트 유틸 절과 예제 9·예제 12입니다.

핵심 개념 (1) 도구 선택 기준표

이 표가 이 도메인의 전부입니다. 다른 것을 못 외웠어도 이건 외우세요.

Kafka 테스트 도구 — 무엇을 검증하는가, 브로커가 필요한가
도구 분류 브로커 무엇을 검증하는가 속도
TopologyTestDriver
kafka-streams-test-utils
단위 불필요 Streams 토폴로지 로직 — 입력을 넣고 출력·상태 저장소를 확인 매우 빠름 (밀리초)
MockProducer
kafka-clients
단위 불필요 "이 코드가 어떤 레코드를 보냈는가" — 전송 목록을 검사 매우 빠름
MockConsumer
kafka-clients
단위 불필요 "레코드가 주어졌을 때 처리 로직이 맞는가" — 레코드를 주입 매우 빠름
Testcontainers
(Kafka 컨테이너)
통합 필요 (컨테이너) 실제 브로커와의 상호작용 — 직렬화, 리밸런스, 트랜잭션, ACL 느림 (수 초~수십 초)
EmbeddedKafkaCluster
(Streams 통합 테스트 유틸)
통합 필요 (JVM 내부) 같은 JVM에서 브로커를 띄워 통합 검증 보통
Spring @EmbeddedKafka
spring-kafka-test
통합 필요 (JVM 내부) Spring 컨텍스트와 함께 리스너·템플릿 검증 보통
kafka-producer-perf-test.sh 성능 필요 쓰기 처리량과 지연 분포
kafka-consumer-perf-test.sh 성능 필요 읽기 처리량
kafka-e2e-latency.sh 성능 필요 end-to-end 지연 (쓰기 → 읽기 왕복)
kafka-verifiable-producer.sh ·
kafka-verifiable-consumer.sh
정합성 필요 유실·중복 없이 전달되었는지 (장애 주입 테스트에 씀)
Datagen 커넥터 · ksqlDB datagen 데이터 생성 필요 스키마 기반 목 데이터 대량 생성

핵심 개념 (2) TopologyTestDriver

kafka-streams-test-utils 아티팩트에 들어 있습니다. Streams 런타임을 시뮬레이션해서, 실제 브로커·컨슈머·프로듀서 없이 토폴로지를 그대로 통과시켜 봅니다.

pom.xml — 테스트 스코프로 추가
<dependency>
  <groupId>org.apache.kafka</groupId>
  <artifactId>kafka-streams-test-utils</artifactId>
  <version>4.3.0</version>
  <scope>test</scope>
</dependency>
WordCountTest.java — 전체 흐름
// 1. 토폴로지를 만든다 (프로덕션 코드와 같은 빌더를 쓴다)
StreamsBuilder builder = new StreamsBuilder();
builder.stream("input-topic")
       .filter((k, v) -> v != null)
       .to("output-topic");
Topology topology = builder.build();

// 2. 테스트 드라이버 생성
TopologyTestDriver testDriver = new TopologyTestDriver(topology);

// 3. 입력 토픽 — 직렬화기(serializer)를 준다
TestInputTopic<String, Long> inputTopic =
    testDriver.createInputTopic("input-topic",
        Serdes.String().serializer(), Serdes.Long().serializer());

// 4. 출력 토픽 — 역직렬화기(deserializer)를 준다
TestOutputTopic<String, Long> outputTopic =
    testDriver.createOutputTopic("output-topic",
        Serdes.String().deserializer(), Serdes.Long().deserializer());

// 5. 데이터를 넣고 결과를 검증
inputTopic.pipeInput("key", 42L);
assertThat(outputTopic.readKeyValue(), equalTo(new KeyValue<>("key", 42L)));

// 6. 상태 저장소를 직접 열어 확인할 수 있다
KeyValueStore<String, Long> store = testDriver.getKeyValueStore("counts-store");
assertThat(store.get("key"), equalTo(1L));

// 7. 벽시계 시각을 전진시켜 punctuation 을 발동시킨다
testDriver.advanceWallClockTime(Duration.ofSeconds(20));

// 8. 반드시 닫는다 (리소스 해제)
testDriver.close();

핵심 개념 (3) MockProducer · MockConsumer

둘 다 kafka-clients에 포함되어 있고 각각 Producer / Consumer 인터페이스의 구현체입니다. 따라서 프로덕션 코드가 인터페이스에 의존해야 갈아끼울 수 있습니다.

MockProducer — "무엇을 보냈는가"를 검사
// autoComplete=true 면 send() 가 즉시 성공 처리된다
MockProducer<String, String> producer =
    new MockProducer<>(true, new StringSerializer(), new StringSerializer());

OrderPublisher publisher = new OrderPublisher(producer);   // 인터페이스 주입
publisher.publish(new Order("order-42", 19900));

List<ProducerRecord<String, String>> sent = producer.history();
assertThat(sent, hasSize(1));
assertThat(sent.get(0).topic(), equalTo("orders"));
assertThat(sent.get(0).key(),   equalTo("order-42"));   // 키를 지정했는지 검증

// 실패 경로도 테스트할 수 있다
MockProducer<String, String> failing =
    new MockProducer<>(false, new StringSerializer(), new StringSerializer());
// … publish 호출 후
failing.errorNext(new org.apache.kafka.common.errors.TimeoutException("boom"));
MockConsumer — 레코드를 주입해 처리 로직을 검증
MockConsumer<String, String> consumer =
    new MockConsumer<>(AutoOffsetResetStrategy.EARLIEST.name());

TopicPartition tp = new TopicPartition("orders", 0);
consumer.assign(List.of(tp));                       // assign 으로 파티션 고정
consumer.updateBeginningOffsets(Map.of(tp, 0L));    // 시작 오프셋을 알려 줘야 한다

consumer.addRecord(new ConsumerRecord<>("orders", 0, 0L, "order-42", "CREATED"));
consumer.addRecord(new ConsumerRecord<>("orders", 0, 1L, "order-42", "PAID"));

new OrderProcessor(consumer, sink).pollOnce();
assertThat(sink.processed(), hasSize(2));

핵심 개념 (4) 통합 테스트

통합 테스트 선택지
방법특징언제 쓰나
Testcontainers Docker로 실제 브로커 이미지를 띄움. 프로덕션과 가장 가까움. Schema Registry·Connect도 함께 띄울 수 있음 직렬화·트랜잭션·보안·리밸런스까지 확인해야 할 때
EmbeddedKafkaCluster 같은 JVM에서 브로커를 기동. Docker 불필요. Apache Kafka의 Streams 통합 테스트에서 쓰는 유틸리티 Docker를 쓸 수 없는 CI, Streams 통합 테스트
Spring @EmbeddedKafka spring-kafka-test가 제공. 테스트 클래스에 애노테이션 하나로 브로커 기동 Spring Boot 애플리케이션의 리스너·템플릿 검증

테스트 가능한 설계 원칙

핵심 개념 (5) 성능 측정 스크립트

프로듀서 성능 측정
bin/kafka-producer-perf-test.sh \
  --topic perf-test \
  --num-records 1000000 \
  --record-size 1024 \
  --throughput -1 \
  --bootstrap-server localhost:9092 \
  --command-property acks=all \
  --command-property linger.ms=5
컨슈머 성능 측정
bin/kafka-consumer-perf-test.sh \
  --topic perf-test \
  --num-records 1000000 \
  --bootstrap-server localhost:9092 \
  --group perf-consumer
주요 옵션과 출력 해석
항목의미
--num-records보낼(읽을) 레코드 수
--record-size레코드 크기(바이트). --payload-file과 배타적
--throughput초당 레코드 상한. -1이면 제한 없음(최대 처리량 측정)
--command-property클라이언트 설정을 커맨드라인에 직접 지정
--command-config설정 파일로 지정
출력 records/sec, MB/sec처리량
출력 avg latency, max latency평균·최대 지연 (프로듀서)
출력 95th, 99th, 99.9th지연 백분위. 평균보다 이게 중요합니다

목 데이터 생성

반드시 외워야 할 이름과 값

클래스·아티팩트·스크립트 이름 (matching 유형 대비)
이름어디에 있는가용도
TopologyTestDriverkafka-streams-test-utils브로커 없이 토폴로지 실행
TestInputTopickafka-streams-test-utilspipeInput()으로 입력 주입 (serializer 필요)
TestOutputTopickafka-streams-test-utilsreadKeyValue()로 출력 검증 (deserializer 필요)
MockProducerkafka-clientshistory()로 전송 레코드 확인
MockConsumerkafka-clientsaddRecord()로 레코드 주입
EmbeddedKafkaClusterStreams 통합 테스트 유틸JVM 내부 브로커
@EmbeddedKafkaspring-kafka-testSpring 통합 테스트
kafka-producer-perf-test.shKafka bin/쓰기 성능
kafka-consumer-perf-test.shKafka bin/읽기 성능
kafka-e2e-latency.shKafka bin/end-to-end 지연
kafka-verifiable-producer.shKafka bin/정합성 검증
kafka-streams-application-reset.shKafka bin/Streams 앱의 오프셋·내부 토픽 초기화 (테스트 반복 실행에 필수)
transaction.state.log.replication.factor브로커 · 기본 3테스트 환경에서 1로
transaction.state.log.min.isr브로커 · 기본 2테스트 환경에서 1로

자주 나오는 함정

함정 1 — TopologyTestDriver vs Testcontainers

둘 다 "Streams를 테스트한다"지만 층위가 다릅니다
관점TopologyTestDriverTestcontainers
무엇인가토폴로지 시뮬레이터실제 브로커를 띄우는 컨테이너 도구
브로커없음있음 (Docker)
언제 발동pipeInput() 호출 시 동기 처리실제 비동기 처리 — 대기(await)가 필요
검증 불가리밸런스, 다중 인스턴스, 브로커 장애(대부분 가능하나 느림)
출제 형태"CI에서 매 커밋마다 돌릴 토폴로지 테스트는?" → TopologyTestDriver

함정 2 — MockProducer vs TopologyTestDriver

둘 다 브로커가 없지만 대상이 다릅니다
관점MockProducer/MockConsumerTopologyTestDriver
대상직접 작성한 프로듀서·컨슈머 코드Streams 토폴로지
아티팩트kafka-clientskafka-streams-test-utils
상태 저장소 검증불가getKeyValueStore()
시간 제어없음advanceWallClockTime()
출제 형태"프로듀서가 키를 지정했는지 검증하려면?" → MockProducer.history()

함정 3 — serializer vs deserializer 방향

입력에는 직렬화기, 출력에는 역직렬화기
관점createInputTopiccreateOutputTopic
필요한 것serializer (Serdes.String().serializer())deserializer
이유Java 객체 → 바이트 (프로듀서 역할)바이트 → Java 객체 (컨슈머 역할)
주요 메서드pipeInput, pipeKeyValueListreadValue, readKeyValue, readRecordsToList
혼동 시 결과컴파일 오류 또는 런타임 ClassCastException
출제 형태코드 스니펫에서 바꿔치기한 인자를 찾는 문항

함정 4 — --throughput -1 vs 값 지정

측정 목적에 따라 반대로 씁니다
관점--throughput -1--throughput 10000
측정 목적최대 처리량목표 부하에서의 지연
지연 수치의 의미큐 대기가 섞여 해석 불가실제 서비스 지연에 근사
혼동 시 결과무제한 부하에서 잰 지연을 SLA 근거로 쓰면 과대 추정
출제 형태"지연 SLA를 검증하려면 어떤 옵션?" → 목표 값 지정

함정 5 — 테스트 환경에서만 실패하는 EOS

클라이언트 설정이 아니라 브로커 설정 문제입니다
관점프로덕션 (브로커 3대)테스트 (브로커 1대)
transaction.state.log.replication.factor3 (기본값 그대로)1로 낮춰야 함
transaction.state.log.min.isr2 (기본값 그대로)1로 낮춰야 함
증상정상트랜잭션 초기화 또는 토픽 생성 실패
혼동 시 결과클라이언트 설정을 계속 고치며 시간을 낭비합니다
출제 형태"단일 브로커에서 exactly_once_v2가 실패하는 이유는?"

함정 6 — 테스트 반복 실행 시 상태가 남는 문제

Streams 통합 테스트를 두 번 돌리면 결과가 달라지는 이유
관점남는 것지우는 방법
컨슈머 오프셋__consumer_offsetsapplication.id 그룹kafka-streams-application-reset.sh
로컬 상태 저장소state.dir 아래 RocksDB 파일KafkaStreams.cleanUp() 또는 디렉터리 삭제
내부 토픽changelog · repartition 토픽reset 도구가 repartition 토픽을 정리
혼동 시 결과"로컬에서는 통과하는데 CI에서 실패" 또는 반대 현상
출제 형태"Streams 앱을 처음부터 다시 처리시키는 도구는?"

코드 읽기 문제 대비

스니펫 1 — 이 테스트는 무엇을 놓치는가

OrderPublisherTest.java
MockProducer<String, Order> producer =
    new MockProducer<>(true, new StringSerializer(), new OrderAvroSerializer());

new OrderPublisher(producer).publish(new Order("order-42", 19900));
assertThat(producer.history(), hasSize(1));

스니펫 2 — 이 테스트가 실패하는 이유

WindowedCountTest.java
// 5분 tumbling 윈도우로 집계하는 토폴로지
TopologyTestDriver driver = new TopologyTestDriver(topology);
TestInputTopic<String, String> in = driver.createInputTopic("events",
    Serdes.String().serializer(), Serdes.String().serializer());

in.pipeInput("user-1", "click");
in.pipeInput("user-1", "click");
in.pipeInput("user-1", "click");

// 윈도우 집계 결과를 기대하지만 아무것도 읽히지 않는다
assertThat(outputTopic.readKeyValue().value, equalTo(3L));

스니펫 3 — 이 성능 측정의 문제

측정 명령과 결과
$ bin/kafka-producer-perf-test.sh --topic t --num-records 5000000 \
    --record-size 1024 --throughput -1 \
    --bootstrap-server localhost:9092 --command-property acks=all

5000000 records sent, 412300 records/sec (402.6 MB/sec),
  1842.11 ms avg latency, 5231.00 ms max latency,
  1701 ms 50th, 4102 ms 95th, 4980 ms 99th, 5180 ms 99.9th

도메인 미니 퀴즈

공식 문서 출처