참고 자료: 이 글은 아래 강의를 들으며 정리하고, 제 나름의 해석을 덧붙인 내용입니다.
https://www.inflearn.com/course/designing-a-server-s
4편에서 응답 속도는 확보했다. 대신 창 하나가 열렸다. - 메시지가 JVM 힙에만 존재해서, 200 OK를 받은 직후 애플리케이션이 죽으면 그 발급 기록은 영원히 사라졌다. 재시도도 불가능했다. - Redis는 이미 그 사람을 발급 완료자로 기록해 뒀기 때문에, 다시 요청해도 "이미 발급됨"으로 거정 당한다.
이번 편에서는 그 큐를 프로세스 밖, Kafka 브로커로 옮긴다. 그런데 "밖으로 옮겼다" 가 "유실 창이 닫혔다" 와 같은 말인지는 따로 확인해야 한다. 실체 코드를 열어보면, 프로듀셔가 브로커에 메시지를 건네는 방식 때문에 창이 훨씬 좁아지긴 해도 완전히 닫히지는 않았다는 걸 확인할 수 있다.
스택: Kotlin · Spring Boot · Spring Data JPA · MySQL 8 · Redis 8 · k6 · Docker Compose
이 글에서 얻어갈 것들
- "큐를 외부화한다"가 실제로 무엇을 사는 것인지. 메시지가 JVM 힙이 아니라 별도 서버의 디스크에 남는다는 것은, 애플리케이션이 죽어도 "당첨됐다"는 사실이 살아남는다는 뜻이다. 그 대신 브로커라는 새로운 운영 대상 하나가 시스템에 추가된다.
- 미들웨어를 도입해도 유실 창이 자동으로 완전히 닫히지는 않는다는 것. 실제 코드를 열어 보면 KafkaTemplate.send()는 브로커의 ack을 기다리지 않고 바로 반환한다. "프로세스 밖으로 옮겼다"와 "유실 창이 사라졌다"는 다른 주장이다.
- 파티션과 컨슈머 그룹이 실제로 하는 일. 왜 유저 아이디를 메시지 키로 쓰는지, 파티션 수와 컨슈머 수를 왜 맞추는지, 그리고 이게 응답 속도와 처리량이라는 서로 다른 두 종류의 성능을 각각 어떻게 개선하는지.
- DLT(Dead Letter Topic)를 가를 때의 기준. 재시도로 흡수할 실패와 격리해서 사람이 봐야 할 실패를 가르지 않으면, 스프링은 그 메시지를 로그 한 줄만 남기고, 조용히 버린다 - DLT가 막아주는 건 파티션이 멈추는 게 아니라, 실패한 메시지가 흔적도 없이 사라지는 것이다.
- 부하 테스트의 경고 메시지도 읽을 줄 알아야 한다는 것. 이번 편의 워밍업은 임계값을 통과하지 못했다. 그런데 원인은 카프카가 아니라 k6의 VU(가상 사용자) 풀이 소진된 것이었고, 이 자리에서 Little's Law로 그걸 확인한다.
0. 4편까지의 상태 - 남은 유실 창
4편에서 "판정"과 "기록"을 시간으로 분리했다. 사용자가 실제로 필요로 하는 정보 - 내가 5000 명 안에 들었는가 - 는 Redis 원자 연산이 실행되는 순간 이미 확정되고, 그 뒤에 벌어지는 DB 기록(카운터 증가 + issuance 행 삽입) 은 응답과 무관하게 뒤에서 처리하면 된다는 발상이었다. 같은 구조로 손으로 만든 인메모리 큐와 스프링 이벤트, 두 가지로 구현했다.
그런데 4편에서 이미 짚어 둔 대로, 둘 다 내구성 프로필이 같았다.
t1 Redis DECR 성공 ← 재고 확정. 이 사용자는 이제 "당첨자"다.
t2 이벤트가 큐에 적재된다 ← LinkedBlockingQueue 또는 ThreadPoolTaskExecutor의 내부 큐, 둘 다 JVM 힙
t3 200 OK 응답 ← 사용자는 여기서 "발급 완료" 화면을 본다
─────────────────────────────────────────────
애플리케이션이 t3~t4 사이에 죽으면?
─────────────────────────────────────────────
t4 워커가 큐에서 이벤트를 꺼내 DB에 INSERT + UPDATE
t3 와 t4 사이 애플리케이션이 재시작,OOM, 강제 종료를 겪으면 큐에 있던 메시지는 그냥 사라진다. 더 나쁜 건, 그 사용자는 재시도도 못 한다는 것이다. - Redis 안에서는 이미 발급 받은 사람이라 재요청은 ALREADY_ISSUED로 막힌다. 재고 한 장이 사라지는 데 더해, 그 재고를 받아야 했던 사람이 영구히 배제된다.
이번 편은 이 t2를 프로세스 밖으로 옮긴다.
1. 큐를 프로세스 밖으로 - 카프카 토픽이라는 외부 저장소
[ 4편 — 스프링 이벤트 ]
발급 API ──► Redis (재고 판정) ──► 200 OK
└──► ApplicationEventPublisher.publishEvent(...)
└──► (같은 JVM 힙 안의 큐) ──► @EventListener가 DB에 기록
[ 5편 — Kafka ]
발급 API ──► Redis (재고 판정) ──► 200 OK
└──► KafkaTemplate.send(topic, key=userId, value=event)
└──► (별도 서버, 디스크) ──► @KafkaListener 워커가 DB에 기록
애플리케이션이 할 일은 두가지다. 토픽으로 연결할 호스트, 포트 설정을 하는 것, 그리고 그 토픽을 구독해 메시지를 가져올 워커(컨슈머)를 만드는 것. 여기서 몇 가지 용어를 정리하고 가자.
- 토픽(topic) — 카프카에서 메시지가 쌓이는 로그(log)의 이름. 지금까지 이 시리즈에서 "큐"라고 불러 온 것과 같은 아키텍처적 역할을 하지만, 정확히는 다르다 — 컨슈머가 읽어도 메시지가 지워지지 않고 보존 기간 동안 그대로 남는다. 바로 아래에서 볼 "오프셋"이 필요한 이유도 여기 있다 — 꺼내면 사라지는 큐라면 오프셋이라는 개념 자체가 필요 없다.
- 오프셋(offset) — 컨슈머가 지금까지 몇 번째 메시지까지 소비했는지 가리키는 위치. 메시지를 처리하고 완료할 때마다 1씩 증가한다.
- 컨슈머 그룹(consumer group) — 이 오프셋이 어디까지 진행됐는지를 카프카 서버가 그룹 단위로 기억한다. 같은 그룹에 속한 컨슈머들이 한 토픽의 파티션들을 나눠 갖는다.
1.1 왜 유저 아이디를 메시지 키로 쓰는가 — 파티션과 병렬 컨슈머
토픽 안에는 파티션이 여러 개 있을 수 있다. 이번 구현은 3개로 정했다. 메시지가 발행되면 어떤 키를 기준으로 각 파티션에 고르게 분배되는데, 여기서는 유저 아이디를 키로 쓴다.
그리고 파티션 하나에 컨슈머(워커) 하나가 1:1로 붙는다. 파티션이 3개면 워커도 3개, concurrency = 3으로 맞춘다.
Producer ──► Topic (issuance.requested)
├── Partition 0 ──► Consumer 0
├── Partition 1 ──► Consumer 1
└── Partition 2 ──► Consumer 2
이 구조가 사는 성능은 두 가지고, 서로 다른 자리에서 일어난다는 걸 구분해야 한다.
- 응답 속도. Coupon API는 메시지를 토픽으로 보내고 바로 응답한다. 이건 파티션 개수와 무관하다 — 프로듀서가 브로커 하나에 메시지를 건네는 일이지, 컨슈머 쪽 병렬성과는 별개다.
- 전체 처리량. 발급 기록을 3개의 워커가 병렬로 처리하면서, 전체 발급 수량을 기록하는 속도 자체가 늘어난다. 이건 순전히 컨슈머 쪽 병렬화의 효과다.
같은 구조를 스프링 이벤트나 인메모리 큐로도 만들 수는 있다. 하지만 유저 아이디로 고르게 분산시켜 각각의 워커 큐로 라우팅하는 로직을 직접 짜야 한다 — 해시 함수를 고르고, 큐를 N개 만들고, 그 사이의 부하 불균형을 다뤄야 한다. 이미 구현이 잘 되어 있고 검증된 미들웨어를 쓰는 게 여기서는 합리적인 선택이다.
이 지점은 4편의 결론과 살짝 결이 다르다는 걸 짚어 둘 만하다. 4편에서는 "손으로 만들어 봐야 스프링이 뭘 대신해 주는지 알 수 있다"고 했다. 이번에도 같은 논리를 적용하면 파티셔닝 로직도 직접 짜 봐야 할 것 같지만, 여기서는 굳이 그러지 않는다. 직접 만들어서 배우는 것과, 검증된 도구에 맡기는 것 중 무엇을 고를지는 그 문제의 복잡도와 실패했을 때의 대가에 달려 있다. 스레드 하나짜리 인메모리 큐를 만드는 것과, 여러 서버에 걸쳐 안전하게 동작하는 분산 해시 라우팅을 만드는 것은 난이도가 다른 급의 문제다.
2. 실패한 메시지를 버리지 않는 법 - DLT
카프카를 쓰면 실패한 메시지를 어떻게 다룰지 정책을 정해야 한다. 컨슈머가 처리하다 실패하면 일단 재처리하는데, 재처리를 해도 계속 실패하는 경우가 있다. 이런 메시지를 위해 별도 토픽 — Dead Letter Topic(DLT) — 을 하나 더 두고, 거기에 옮겨 나중에 사람이 보도록 한다. (큐 진영에서는 DLQ라고 더 많이 부르지만, 카프카는 큐 단위가 토픽이라 DLT라는 표현을 쓴다.)
어떤 실패를 DLT로 보내야 하는가?
| DLT로 보내는 경우 | 이유 |
| DB 연결이 계속 타임아웃 | 일시적 장애가 아니라 지속되는 장애 |
| 비즈니스 로직 버그로 인한 NPE | 발생하면 안 되는 에러 — 사람이 봐야 한다 |
| 메시지 포맷 오류 | 프로듀서·컨슈머가 서로 정한 포맷과 어긋남 |
| DB 정합성 위반 (외래키 등) | 예상치 못한 데이터 상태 |
| 알 수 없는 예외 전파 | NPE와 같은 맥락 |
| N회 재시도 후에도 계속 실패 | 정해진 시간 안에 처리할 수 없다고 판단 |
반대로 DLT로 보내지 않는 경우도 명확히 있다.
| DLT로 보내지 않는 경우 | 이유 |
| 유니크 제약 위반을 catch로 잡은 경우 | 첫 성공 요청 이후의 중복 메시지 — 정상 흐름 |
| 정상 처리 후 반환 | 성공 |
| 오프셋 커밋 | 처리가 끝났다는 뜻이므로 그 자체가 성공 사례 |
여기서 한 가지 더 짚어야 할 위험이 있다. DLT라는 안전망 없이 실패한 메시지를 계속 재시도만 하도록 놔두면 어떻게 되는가. 오프셋은 처리가 끝나야 올라간다. 즉 같은 메시지가 계속 실패하는 동안 그 파티션의 오프셋은 전혀 전진하지 못하고, 그 뒤에 쌓인 나머지 메시지들도 전부 대기 상태로 묶인다. 파티션 하나는 유저 아이디 기반으로 분배되므로, 하필 그 실패 메시지와 같은 파티션에 걸린 사용자들만 골라서 발급 기록이 멈추는 셈이다. 파이프라인 전체는 "정상 동작 중"으로 보이는데, 특정 사용자군만 조용히 지연되는 상황이 생길 수 있다. N회 재시도 뒤 DLT로 격리하는 정책이 없으면, 메시지 하나의 실패가 그 파티션 전체를 사실상 멈춰 세운다.
3 실습 - 스프링 이벤트를 카프카로 교체하기
build.gradle.kts에 의존성 하나를 추가한다
implementation("org.springframework.kafka:spring-kafka")
3.1 docker-compose.yml - Kafka 3.8 (KRaft, 단일브로커)
kafka:
# Kafka 3.8부터 KRaft 모드가 기본이라 ZooKeeper가 없다.
image: apache/kafka:3.8.0
environment:
KAFKA_NODE_ID: 1
KAFKA_PROCESS_ROLES: broker,controller
KAFKA_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093,HOST://:9094
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092,HOST://localhost:9094
KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka:9093
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,CONTROLLER:PLAINTEXT,HOST:PLAINTEXT
KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
KAFKA_NUM_PARTITIONS: 3
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
KAFKA_AUTO_CREATE_TOPICS_ENABLE: "true"
ports:
- "9094:9094"
healthcheck:
test: ["CMD-SHELL", "/opt/kafka/bin/kafka-broker-api-versions.sh --bootstrap-server localhost:9092"]
interval: 10s
timeout: 5s
retries: 20
몇 가지 값은 반드시 설명이 필요하다.
노드가 1개인 이유. 로컬 개발 환경에서 리소스를 아끼기 위한 단일 브로커 구성이다. 운영에서는 보통 3개 또는 5개로 구성한다. 이 차이는 5절에서 다시 중요해진다 — 노드가 하나라는 건 곧 복제본도 하나뿐이라는 뜻이고, 이번 구성의 카프카 자체가 단일 장애점(SPOF)이라는 뜻이다.
리스너가 세 개인 이유. 9092는 같은 도커 네트워크 안의 다른 컨테이너(애플리케이션)가 접속하는 내부 주소다. 9093은 컨트롤러 전용 내부 통신용이다. 9094는 로컬 맥북에서 직접 카프카에 프로듀서/컨슈머로 접속해 테스트해 보고 싶을 때 쓰는 외부용 포트다. 실제 application.yaml을 보면 spring.kafka.bootstrap-servers의 기본값이 localhost:9094로 잡혀 있어서, 애플리케이션을 도커 밖(로컬)에서 직접 띄워도 별도 설정 없이 카프카에 붙을 수 있게 해 뒀다.
AUTO_CREATE_TOPICS_ENABLE: true. 애플리케이션이 어떤 토픽에 연결하면 그 토픽이 카프카에 없을 경우 자동으로 만들어진다. 개발 편의를 위한 설정이고, 운영 카프카 클러스터에서는 보통 이렇게 두지 않는다. 오타나 실수로 생성된 토픽이 그대로 살아남을 수 있기 때문이다.
REPLICATION_FACTOR: 1. 카프카가 여러 노드를 두는 핵심 이유 중 하나가 데이터를 여러 노드에 복제해 내구성을 높이는 것이다. 3중 복제라면 노드 하나가 죽어도 나머지 두 곳에 데이터가 남는다. 이번 구성은 노드가 하나뿐이라 복제 자체가 의미가 없고, 그래서 1로 뒀다.
3.2 프로듀서 - IssuanceRequestedProducer
@Component
class IssuanceRequestedProducer(
private val kafkaTemplate: KafkaTemplate<String, Any>,
) {
fun publish(event: IssuanceRequested) {
kafkaTemplate.send(IssuanceTopics.REQUESTED, event.userId.toString(), event)
}
}
키는 유저 아이디를 문자열로, 값은 IssuanceRequested 객체를 그대로 넘기고 직렬화는 설정(3.5절)에서 JSON으로 맡긴다. CouponService.issue()를 실제로 열어 보면, 4편에서 eventPublisher.publishEvent(...)였던 자리가 정확히 issuanceRequestProducer.publish(event) 한 줄로 바뀌어 있을 뿐, 그 앞뒤 순서(couponIssuer.tryIssue(...) 먼저, 그다음 발행)는 그대로다.
토픽 이름과 컨슈머 그룹 이름은 한 오브젝트에 상수로 모아 둔다.
object IssuanceTopics {
const val REQUESTED = "issuance.requested"
const val REQUESTED_DLT = "issuance.requested.DLT"
const val CONSUMER_GROUP = "issuance-worker"
}
3.3 컨슈머 - IssuanceWorker
@Component
class IssuanceWorker(
private val writer: IssuanceWriter,
) {
@KafkaListener(
topics = [IssuanceTopics.REQUESTED],
groupId = IssuanceTopics.CONSUMER_GROUP,
concurrency = "3",
)
fun consume(event: IssuanceRequested) {
writer.write(event)
}
}
writer.wirte(event) 는 4편에서 만들어 둔 IssuanceWorker 그대로다. - 유니크 제약 위반을 조용히 흡수하는 멱등 처리 포함해서. 4편에서는 지금 당장은 한번도 발동하지 않지만 카푸카로 넘어가면 실제로 일 할 차례가 온다. 카프카는 최소 한 번(at-least-once) 전달만 보장하고 정확히 한 번(exactly-once)을 보장하지 않으므로, 리벨런싱이나 재시도 과정에서 같은 메시지가 두 번 소비될 수 있다.
DataIntegrityViolationException을 흡수하는 그 코드가 바로 그 중복을 막는 유일한 방어선이다.
concurrency = 3 은 파티션(3)에 맞춘 값이다. 파티션보다 컨슈머가 많으면 남은 컨슈머는 그냥 논다 - 파티션과 컨슈머는 1:N 이 아니라 최대 1:1 까지만 유효하다.
3.4 토픽 자동 생성 — KafkaTopicConfig
@Configuration
class KafkaTopicConfig {
@Bean
fun issuanceRequestedTopic(): NewTopic =
TopicBuilder.name(IssuanceTopics.REQUESTED).partitions(3).replicas(1).build()
@Bean
fun issuanceRequestedDltTopic(): NewTopic =
TopicBuilder.name(IssuanceTopics.REQUESTED_DLT).partitions(3).replicas(1).build()
}
애플리케이션이 시작할 때 이 빈들이 등록되어 있으면, 카프카에 해당 토픽이 없을 경우 새로 만들어준다.
3.5 프로듀셔, 컨슈머 설정 - KafkaConfig
ProducerConfig.ACKS_CONFIG to "1"
메시지를 보냈을 때 성공으로 칠 기준이다. acks = 1 은 리더 브로커 하나에만 도달해도 성공으로 본다. acks=all(또는 -1) 로 두면 모든 복제본에 반영된 걸 확인하고서야 성공으로 치므로 내구성은 더 높아지지만 응답이 느려진다. 이번 구성은 브로커가 하나뿐이라 acks=1과 acks=all이 사실상 같은 의미다. 여러 브로커로 운영한다면 이 값 하나가 지연 시간과 내구성 사이에 실제 트레이드 오프가 된다.
직렬화 쪽에서 실제 코드를 보고 알게 된 게 하나 있다. 이번 저장소는 spring-boot-starter-webmvc 계열의 최신 스프링 부트를 쓰고 있어서, JSON 직렬화기가 익숙한 org.springframework.kafka.support.serializer.JsonSerializer/JsonDeserializer가 아니라 JacksonJsonSerializer / JacksonJsonDeserializer 다.
내부적으로 tools.jackson.databind.json.JsonMapper(Jackson 3 계열의 새 패키지 네임스페이스)를 쓴다. 온라인에 떠도는 예전 스프링 카프카 예제 코드를 그대로 복사하면 이 부분에서 컴파일이 안 될 수 있다 — 클래스 이름과 패키지가 바뀌었기 때문이다. 새로 시작하는 프로젝트일수록 사용 중인 스프링 부트·스프링 카프카 버전이 어떤 직렬화기를 기본으로 미는지 먼저 확인하는 습관이 필요하다.
컨슈머 설정에서는 두 가지가 중요하다.
컨슈머 설정에서는 두 가지가 중요하다.
ConsumerConfig.GROUP_ID_CONFIG to IssuanceTopics.CONSUMER_GROUP
ConsumerConfig.AUTO_OFFSET_RESET_CONFIG to "earliest"
그리고 역직렬화기는 이렇게 감싼다.
val jsonDelegate = JacksonJsonDeserializer(IssuanceRequested::class.java, jsonMapper).apply {
addTrustedPackages("com.apiece.coupon.infrastructure.messaging")
}
val valueDeserializer = ErrorHandlingDeserializer(jsonDelegate)
auto-offset-reset = earliest는 컨슈머가 토픽에 처음 붙었을 때 오프셋 기록이 없으면 맨 처음부터 다시 읽으라는 의미다. 이번 실습에서는 부하테스트와 개발 편의을 위해 이렇게 뒀지만, 운영에서 이 값을 그대로 쓰면 재시작 시 메시지를 중복으로 다시 읽을 위험이 있다. 보통은 latest 나 none을 쓴다.
addTrustedPackages(...)는 보안과 직결된 설정이다. 역질렬화는 임의이 바이트를 객체로 바꾸는 과정이라, 아무 클래스나 허용하면 악의적으로 조작된 메시지가 원치 않는 코드를 실행시키는 통로가 될 수 있다. 미리 허용한 패키지의 하위 타입만 역질렬화를 허용하는 것이 최소한의 방어선이다.
kafkaListenerContainerFactory 빈에는 ackMode나 ContainerProperties 를 건드리는 코드가 없다. 즉 BATCH를 명시적으로 선택한게 아니라 스프링 카프카의 기본값을 그대로 쓰고 있는 상태다. 지금 워크로드(파티션마다 메시지를 한 건씩 받아 한건씩 커밋)에서는 기본값으로 충분하지만, "의도적으로 고른 기본값"과 "그냥 안 건드린 기본값"은 코드만 봐서는 구분이 안 된다. 나중에 이 설정을 바꿀 일이 생기면, 지금 이 값이 의도된 것인지부터 확인해야 한다.
3.6 DLT 라우터 - KafkaErrorHandlerConfig
@Configuration
class KafkaErrorHandlerConfig {
@Bean
fun errorHandler(template: KafkaTemplate<String, Any>): DefaultErrorHandler {
val recoverer = DeadLetterPublishingRecoverer(template) { record, _ ->
TopicPartition("${record.topic()}.DLT", record.partition())
}
return DefaultErrorHandler(recoverer, FixedBackOff(1_000, 3))
}
}
정책은 이렇다. 워커가 메시지 처리에 실패하면 1초 쉬고 재시도한다. 그래도 실패하면 또 1초 쉬고 재시도, 최대 3번까지다. 3번 모두 실패하면 그제서야 issuance.requested.DLT로 - 원본과 같은 파티션 번호를 유지한 채 - 옮겨진다. 그 뒤 사람이 DLT 토픽을 보고 원인을 조사하고 재처리한다.
3.7 프로듀서의 send()는 브로커의 ack을 기다리지 않는다.
3.5절의 acks=1과는 다른 층위의 이야기라는 걸 먼저 구분해 두는 게 좋다 — acks는 브로커가 무엇을 확인해야 성공으로 칠지를 정하는 설정이고, 지금 볼 것은 애플리케이션이 그 확인을 실제로 기다리는지의 문제다. acks=all로 아무리 강하게 잡아도, 호출한 쪽이 결과를 기다리지 않으면 그 설정은 브로커 내부에서만 의미가 있다.
IssuanceRequestedProducer.publish()가 호출하는 kafkaTemplate.send(...)는 CompletableFuture 를 즉시 반환하는 비동기 호출이다. 실제 코드를 보면 이 반환값을 어디서도 받지 않는다. - .get()을 걸어 기다리지도, 콜백을 등록하지도 않는다. 호출한 스레드는 이 한줄을 지나 바로 issue()의 나머지 (응답 객체 생성) 로 넘어간다.
즉 issue()가 publish()를 호출하고 바로 다음 줄로 넘어가 200 OK 를 반환하는 시점에, 메시지는 아직 프로듀서의 내부 배치 버퍼에 있을 뿐 브로커 디스크에 실제로 기록됐다는 보장이 없다. 이 사이에도 창이 하나 있다는 뜻이다.
t1 Redis DECR 성공
t2 kafkaTemplate.send() 호출 ← 프로듀서 내부 버퍼에 적재, 아직 브로커로 안 갔을 수 있음
t3 200 OK 응답
─────────────────────────────────────────
t3~t4 사이 애플리케이션이 죽으면, 그리고 t2의 메시지가 아직 브로커에 도달 전이었다면?
─────────────────────────────────────────
t4 브로커가 메시지를 실제로 받아 디스크에 기록
t5 워커가 소비해 DB에 기록
4편의 t2~t4 창(JVM 힙에서 워커가 소비하기까지)과 비교하면 이 창은 훨씬 좁다. 프로듀서 내부 버퍼가 네트워크로 플러시되는 시간은 보통 수 밀리초 단위다. 반면 4편의 창은 "워커가 처리를 마칠 때까지" 전체 구간이었다. 그러니 이번 구조가 확실히 유실 창을 훨씬 좁힌 건 맞다. 하지만 "완전히 닫았다"는 다른 주장이고, 코드가 그걸 뒷받침하지 않는다.
완전히 닫으려면 send()의 반환값을 동기적으로 기다리거나, 최소한 실패 시 재시도·보상 로직이 딸린 콜백을 붙여야 한다. 그런데 그렇게 하면 응답 경로에 다시 네트워크 왕복이 하나 들어오는 셈이라, 4편에서 그토록 걷어내려 했던 지연이 다른 이름으로 돌아온다. 여기서도 같은 질문이 반복된다 — 무엇을 포기했는지 알고 있는가. 이번 코드의 선택은 "아주 좁은 유실 창"과 "응답 속도"를 맞바꾼 것이지, 유실 창을 아예 없앤 게 아니다.
4. 측정 — 워밍업이 실패한 이유는 카프카가 아니라 VU풀이었다.
같은 조건(재고 5,000 장, 초당 5,000 건 도착률, 30초) 로 워밍업 1회, 본 측정 1회를 순서대로 돌렸다.
4.1 워밍업 - 임계값 실패 - 그런데 원인이 다르다.
WARN[0003] Insufficient VUs, reached 5000 active VUs and cannot initialize more
issue_latency: ✗ 'p(99)<500' p(99)=1.64s
avg=198.59ms p(90)=918.37ms p(95)=1.28s
dropped_iterations: 7415
vus_max: 5000 (min=2600, max=5000)
임계값(p99 < 500ms)이 깨졌다. 그런데 이걸 "카프카를 붙였더니 느려졌다"로 읽으면 틀린다. 로그 첫 줄이 원인을 이미 말해주고 있다 — k6가 VU(가상 사용자)를 5,000개까지 다 썼는데도 더 필요해서 요청을 만들지 못했다.
이게 왜 일어나는지는 Little's Law로 설명된다. 초당 5,000건의 도착률을 유지하려면, 필요한 동시 실행 수(VU)는 대략 도착률 × 평균 처리 시간이다. 워밍업의 평균 지연은 198.59ms였으니 5,000 × 0.198 ≈ 990개 정도가 평균적으로 필요하다. 문제는 p95가 918ms, p99가 1.26~1.64까지 치솟는 꼬리다 - 꼬리 구간에서는 5,000 * 1.64 = 8,200 개가 순간적으로 필요해진다. k6 스크립트에 설정된 maxVUs : 5000 이 순간 수요를 못 버틴것이다.
왜 이번 편의 워밍업만 이렇게 느렸는가. 4편의 성공 경로는 Redis 원자 연산 + 인메모리 큐 적재로 끝났고, 이건 마이크로초 단위였다. 이번엔 그 자리에 네트워크를 타는 카프카 프로듀서 호출이 들어간다. 게다가 컨슈머 쪽은 이제 컨슈머 그룹이 파티션에 처음 조인하는 비용, 토픽·파티션 메타데이터를 받아오는 비용까지 콜드 스타트로 얹힌다 — run.sh의 주석에 미리 "컨슈머 그룹 조인 비용 흡수"라고 적어 둔 게 바로 이거다.
즉 워밍업의 실패는 시스템이 느리다는 증거가 아니라, 콜드 스타트 비용이 어느 정도인지 보여주는 증거다. 그리고 k6의 VU 상한을 이 콜드 스타트를 감당할 만큼 넉넉히 잡지 않으면, 시스템이 아니라 부하 생성기 자체가 병목이 되어 버린다는 걸 이번 결과가 정확히 보여준다.
4.2 본 측정 — 임계값 통과, 그런데 무엇과 비교할 것인가
issue_latency: ✓ 'p(99)<500' p(99)=51.04ms
avg=4.34ms med=1.72ms p(90)=7.96ms p(95)=19.06ms
{ expected_response:true }: avg=17.97ms min=602µs med=8.63ms max=179.04ms p(90)=41.15ms p(95)=61.22ms
http_req_failed: 96.66% (145000 / 150000)
vus: 13 (min=6, max=189)
워밍업으로 컨슈머 그룹 조인과 커넥션 풀이 이미 끝난 상태에서 다시 돌리니 p99가 51.04ms로 떨어지며 임계값을 여유 있게 통과했다. 발급 5,000장, 정합성 문제 없음.
k6는 http_req_duration에 { expected_response:true }라는 서브 지표를 자동으로 같이 찍는다 — 응답 상태가 성공 범위(2xx~3xx)인 요청만 따로 모은 값이다. 성공 요청만의 평균은 17.97ms, p95는 61.22ms로, 전체를 섞은 issue_latency의 avg 4.34ms와는 꽤 차이가 크다. 전체 15만 건 중 성공은 5,000건(3.3%)뿐이고, 나머지 96.7%는 Redis 게이트에서 1~2ms 만에 끝나는 거절이라 섞인 평균이 성공 쪽 신호를 가려 버리기 때문이다. 3편에서 이미 한 번 짚었던 원칙이 이번에도 그대로 적용된다 — 집계 지표는 실제로 일어난 일의 대부분을 가린다. 이 시스템의 응답성을 말하고 싶다면, "전체 평균 4.34ms"가 아니라 "성공한 사람은 평균 18ms, 느려도 p95가 61ms 안에 끝난다"는 쪽이 훨씬 정직한 문장이다.
5. 이번 편이 산 것, 그리고 아직 안 산 것
| 4편 (스프링 이벤트) | 5편 (카프카) | |
| 메시지가 사는 곳 | 같은 JVM 힙 | 별도 브로커 디스크 |
| 애플리케이션 재시작 시 | 큐에 남은 메시지 전부 유실 | 브로커에 남아 재시작 후 이어서 처리 |
| 프로듀서 전송 | 인메모리 함수 호출 (사실상 즉시 반영) | 네트워크 호출, ack 대기 없이 반환 (3.7절) |
| 남은 유실 창 | 크다 — 워커가 처리를 마칠 때까지 전체 | 훨씬 좁다 — 프로듀서 버퍼가 플러시될 때까지 |
| 브로커/미들웨어 장애 시 | 해당 없음 (미들웨어가 없음) | 이번 구성(단일 노드, 복제 없음)에서는 카프카가 곧 새로운 SPOF |
| 중복 메시지 방어 | 필요는 하지만 발동 안 함 | at-least-once 전달로 실제 발동 가능 — 4편의 유니크 제약 방어가 여기서 일한다 |
정리하면 이번 편은 유실 창을 훨씬 좁히는 대가로 브로커라는 새 운영 대상 하나를 들여왔다. 그 브로커를 단일 노드·복제 없음으로 구성한 한, 애플리케이션 프로세스가 갖고 있던 단일 장애점이 카프카 브로커로 자리만 옮긴 셈이기도 하다. 운영 환경이라면 최소 3개 노드, 복제 팩터 2~3 구성이 이 절의 마지막 줄을 다르게 만든다.
6. 닫으며
2편부터 이어 온 패턴에 한 줄을 더 보탠다.
- 비관적 락은 정확성을 얻고 동시성을 팔았다.
- 조건부 UPDATE는 정확성을 지키며 동시성도 되찾았지만, 거절 트래픽은 여전히 DB가 받았다.
- Redis 원자 연산은 거절을 1밀리초 미만으로 만들었지만, 정확한 N개를 최대 N개로 낮췄다.
- 큐 분리(4편)는 응답 속도를 샀지만, 그 응답이 나간 뒤에도 사실이 남아 있다는 보장은 팔았다.
- 큐를 프로세스 밖으로 옮긴 것(5편)은 그 유실 창을 훨씬 좁혔지만, 완전히 없애지는 못했다. 프로듀서 호출 자체가 비동기인 한, 그리고 브로커가 단일 노드인 한, 아주 작지만 여전히 0이 아닌 창과 새로운 단일 장애점이 남는다.
같은 문장을 다시 적는다
무엇을 포기했는지 알고 있는가
이번 편에서도 답은 "전부 없앴다"가 아니었다. 유실 창을 초 단위에서 밀리초 단위로 줄인 것, 그리고 그 대가로 운영해야 할 시스템이 하나 늘어난 것 — 이 둘을 같이 봐야 이번 편이 정직하게 산 것과 안 산 것을 구분할 수 있다.