[Kafka] 카프카 컨슈머와 스프링 카프카

개요

Kafka 컨슈머는 파티션에서 메시지를 읽어 비즈니스 로직으로 넘기는 주체다. 내부적으로는 단일 스레드에서 도는 poll() 루프 하나를 중심으로 Fetcher, Coordinator, Deserializer가 맞물려 움직인다. 이 구조를 이해해야 처리량을 어디서 올릴 수 있고, 어디서 병목이 생기는지 보인다.

컨슈머 여러 개를 하나의 그룹으로 묶으면 토픽의 파티션을 나눠 가진다. 이때 컨슈머가 죽거나 추가되면 파티션을 재배분하는 리밸런싱이 일어나고, 그 방식에 따라 처리량과 가용성이 갈린다. 리밸런싱을 어떻게 다루느냐가 운영 안정성의 핵심이다.

Spring Kafka는 이 모든 동작을 @KafkaListenerConcurrentMessageListenerContainer로 추상화한다. AckMode로 커밋 시점을 고르고, ErrorHandler와 @RetryableTopic으로 실패를 재시도하거나 격리한다. 그 위에서 처리량을 올리는 세 가지 접근(배치 컨슈머, 비동기 컨슈머, Confluent Parallel Consumer)을 측정 결과와 함께 비교한다.

Consumer 내부 동작

전체 흐름

컨슈머는 단일 스레드에서 동작하는 poll() 루프를 중심으로 움직인다.

  1. Coordinator 통신: 백그라운드 하트비트 스레드가 별도로 동작하고 있고, poll 시점에 리밸런싱 트리거 확인, 자동 오프셋 커밋 등이 함께 처리된다.
  2. Fetcher: 할당된 파티션의 리더 브로커에 Fetch 요청을 보낸다. 미리 가져오기(prefetch)로 대기 시간을 줄인다.
  3. Deserializer: 바이트 배열을 Java 객체로 역직렬화한다.
  4. Interceptor: 메트릭, 로깅, 트레이싱 같은 횡단 관심사를 처리한다.

KafkaConsumer는 thread-safe하지 않다. 한 인스턴스를 여러 스레드에서 동시에 사용하면 안 된다. Spring Kafka의 ConcurrentMessageListenerContainer도 내부적으로는 컨테이너마다 별도의 컨슈머 인스턴스를 두고 각각 단일 스레드로 운영한다.

한 스레드로는 처리량에 한계가 있어 병렬로 소비하고 싶지만, 컨슈머 1개를 여러 스레드가 공유하는 방식은 thread-safe하지 않아 쓸 수 없다. 컨슈머는 자신이 할당받은 파티션의 오프셋과 내부 상태를 단일 스레드 기준으로 관리하기 때문이다. 여러 스레드가 한 인스턴스를 공유하면 같은 메시지가 중복 처리되거나 오프셋 관리가 깨질 수 있다.

@KafkaListener(concurrency = 3)처럼 동시성을 주면 Spring Kafka는 한 컨슈머를 멀티스레드로 돌리는 것이 아니라, 내부에 실행 단위(컨테이너)를 3개 만들고 각 컨테이너마다 KafkaConsumer 인스턴스를 1개씩 둔다. 각 컨슈머는 자기만의 단일 스레드에서 poll() 루프를 돈다. 즉 동시성 3은 KafkaConsumer 3개가 각자 단일 스레드로 도는 구조이며, thread-safe 규칙을 지키면서 병렬성을 얻는다. 이 내부 구조는 뒤의 concurrency의 진짜 의미 절에서 다이어그램으로 다시 정리한다.

단, concurrency 값은 파티션 수를 넘으면 의미가 없다. 컨슈머 3개를 띄워도 파티션이 2개면 컨슈머 1개는 유휴 상태가 되므로, 보통 concurrency를 파티션 수 이하로 맞춘다.

Fetcher 설정

Fetcher의 동작은 다음 파라미터들로 조정한다.

파라미터기본값의미
fetch.min.bytes1응답 최소 크기. 이 값에 도달하기 전까지 브로커가 대기
fetch.max.wait.ms500fetch.min.bytes 미달 시 최대 대기 시간
max.partition.fetch.bytes1048576 (1MB)파티션당 한 번에 가져올 최대 바이트
fetch.max.bytes52428800 (50MB)전체 요청의 최대 바이트
max.poll.records500한 번의 poll()에서 반환할 최대 레코드 수

fetch.min.bytes를 크게 설정하면 적은 메시지를 모아 한 번에 가져오므로 처리량이 올라가지만, 그만큼 첫 메시지의 지연이 늘어난다.

컨슈머 그룹과 리밸런싱

컨슈머 그룹의 기본 규칙

여러 컨슈머가 하나의 그룹으로 묶이면, 그룹 안에서 토픽의 파티션을 나눠 가진다.

  • 한 파티션은 그룹 내 하나의 컨슈머에만 할당된다.
  • 컨슈머 수가 파티션 수보다 많으면 남는 컨슈머는 유휴 상태가 된다.
  • 컨슈머 수가 파티션 수 이하가 되도록 설계해야 한다.
  • 서로 다른 그룹은 같은 데이터를 독립적으로 소비한다.

리밸런싱이란

리밸런싱이란 컨슈머 그룹의 구성이 바뀔 때 그룹 코디네이터(브로커)가 토픽의 파티션을 그룹 내 컨슈머들에게 다시 나눠 주는 과정이다. 다음 상황에서 발생한다.

  • 컨슈머 그룹에 새 컨슈머 가입
  • 컨슈머 그룹에서 컨슈머 탈퇴 (정상 종료 또는 장애)
  • 구독 토픽의 파티션 수 변경
  • session.timeout.ms 내에 하트비트 미수신
  • max.poll.interval.ms 내에 poll() 호출 없음

Eager vs Cooperative 리밸런싱

Eager 리밸런싱은 모든 컨슈머가 가진 파티션을 일단 모두 해제(revoke)하고 다시 할당받는다. 그 동안 그룹 전체의 소비가 중단된다(Stop-the-world).

Cooperative (Incremental) 리밸런싱은 이동이 필요한 파티션만 revoke하고 나머지는 계속 소비한다. Kafka 2.4부터 CooperativeStickyAssignor로 사용 가능하다. 신규 컨슈머가 들어와도 다른 컨슈머는 자기 파티션의 일부만 잠시 내려놓을 뿐 소비가 멈추지 않는다.

sequenceDiagram
    participant C1 as Consumer 1
    participant C2 as Consumer 2
    participant CO as Coordinator
    participant C3 as Consumer 3 (신규)
    Note over C1,C2: 정상 처리 중
    C3->>CO: JoinGroup
    CO->>C1: 일부 파티션만 revoke
    Note over C2: 계속 소비
    CO->>C3: revoke된 파티션 할당
    Note over C1,C3: 부분적 중단만 발생
partition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignor

파티션 할당 전략

전략동작특징
RangeAssignor (기본)토픽 단위로 연속 범위 할당토픽 많을 때 앞쪽 컨슈머에 부하 편중
RoundRobinAssignor모든 파티션을 라운드로빈으로 분배고르게 분배되나 리밸런싱 시 이동 많음
StickyAssignor기존 할당 최대한 유지 + 균등 분배리밸런싱 이동 최소화
CooperativeStickyAssignorStickyAssignor + Cooperative 프로토콜부분적 리밸런싱 + 이동 최소화 (권장)

Static Group Membership

Kubernetes 환경에서 롤링 재시작 시 매번 리밸런싱이 일어나는 문제를 해결하려면 group.instance.id를 설정한다.

group.instance.id=order-consumer-pod-1
session.timeout.ms=120000  # 재시작 시간을 고려해 충분히 크게

같은 ID로 다시 접속하면 코디네이터는 같은 인스턴스가 재시작한 것으로 판단하고 기존 파티션을 그대로 다시 할당한다. session.timeout.ms 안에 재접속해야 한다.

Spring Kafka Consumer

@EnableKafka와 @KafkaListener

Spring Kafka는 @EnableKafka로 활성화된다. 이 어노테이션이 KafkaListenerAnnotationBeanPostProcessor를 등록하고, 이 BPP가 @KafkaListener가 붙은 메서드를 스캔해서 리스너 컨테이너로 변환한다.

@Configuration
@EnableKafka
public class KafkaConsumerConfig {
 
    @Bean
    public ConsumerFactory<String, Object> consumerFactory() {
        return new DefaultKafkaConsumerFactory<>(Map.of(
                ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092",
                ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class,
                ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class,
                ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false,
                ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest",
                JsonDeserializer.TRUSTED_PACKAGES, "com.example.*"));
    }
 
    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, Object> kafkaListenerContainerFactory(
            ConsumerFactory<String, Object> factory) {
        var container = new ConcurrentKafkaListenerContainerFactory<String, Object>();
        container.setConsumerFactory(factory);
        container.setConcurrency(3);
        container.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL);
        return container;
    }
}

concurrency

concurrency는 하나의 애플리케이션 인스턴스에서 몇 개의 컨슈머 스레드(= KafkaConsumer 인스턴스)를 띄울지를 의미한다.

concurrency가 파티션 수보다 크면 자동으로 파티션 수로 줄어든다. 남는 스레드는 idle 상태가 된다. 즉, concurrency로 단일 파티션 안에서 여러 스레드가 동시에 읽도록 만들 수는 없다. 단일 파티션을 병렬 처리하려면 뒤에서 다룰 Confluent Parallel Consumer 같은 다른 접근이 필요하다.

리스너 사용 패턴

@Component
@Slf4j
@RequiredArgsConstructor
public class OrderConsumer {
 
    private final OrderService orderService;
 
    // 단건 처리 + 수동 ack
    @KafkaListener(topics = "orders", groupId = "order-service",
                   containerFactory = "kafkaListenerContainerFactory")
    public void listen(@Payload OrderEvent event,
                       @Header(KafkaHeaders.RECEIVED_PARTITION) int partition,
                       @Header(KafkaHeaders.OFFSET) long offset,
                       Acknowledgment ack) {
        orderService.process(event);
        ack.acknowledge(); // 처리 완료 후 커밋
    }
 
    // 배치 처리
    @KafkaListener(topics = "events", groupId = "batch-processor", batch = "true")
    public void listenBatch(List<ConsumerRecord<String, Event>> records,
                            Acknowledgment ack) {
        eventService.processBatch(records.stream().map(ConsumerRecord::value).toList());
        ack.acknowledge();
    }
 
    // 여러 토픽 처리
    @KafkaListener(topics = {"orders", "payments", "shipments"})
    public void listenMultipleTopics(ConsumerRecord<String, Object> record) {
        switch (record.topic()) {
            case "orders" -> handleOrder(record);
            case "payments" -> handlePayment(record);
            case "shipments" -> handleShipment(record);
        }
    }
}

예외를 메서드 밖으로 던지면 ErrorHandler가 받아 재시도/DLT 처리한다. try-catch로 삼키면 ErrorHandler가 동작하지 않는다.

AckMode 비교

AckMode커밋 시점
RECORD레코드 처리 후 즉시 커밋
BATCH (기본)poll로 받은 배치 전체 처리 후 커밋
TIMEackTime 간격으로 커밋
COUNTackCount 개 처리 후 커밋
COUNT_TIMECOUNT 또는 TIME 중 먼저 도달 시 커밋
MANUALAcknowledgment.acknowledge() 호출 시
MANUAL_IMMEDIATEacknowledge() 즉시 커밋 (배치 안 기다림)

에러 핸들링과 재시도

컨슈머가 메시지를 처리하다 실패하면 어떻게 에러 핸들링을 하고, 재시도를 할 것인지 구성해야한다.

유형회복 가능성대응
외부 시스템 일시 장애자동 재시도로 회복ErrorHandler 재시도
처리 불가한 데이터 (유효성 실패)메시지 가공 필요DLT로 격리
스키마 불일치스키마 버전 맞으면 자동 회복컨슈머 미커밋 후 자동 재시도
순서 위반비즈니스 판단 필요알림 + 수동 개입
리밸런싱 중 처리자동 회복멱등성 보장

DefaultErrorHandler (블로킹 재시도)

@Bean
public ConcurrentKafkaListenerContainerFactory<String, Object> kafkaListenerContainerFactory(
        ConsumerFactory<String, Object> factory, KafkaTemplate<String, Object> kafkaTemplate) {
 
    var container = new ConcurrentKafkaListenerContainerFactory<String, Object>();
    container.setConsumerFactory(factory);
 
    var recoverer = new DeadLetterPublishingRecoverer(kafkaTemplate,
            (record, ex) -> new TopicPartition(record.topic() + ".DLT", record.partition()));
 
    var backOff = new ExponentialBackOff();
    backOff.setInitialInterval(1000L);
    backOff.setMultiplier(2.0);
    backOff.setMaxInterval(10000L);
    backOff.setMaxElapsedTime(30000L);
 
    var errorHandler = new DefaultErrorHandler(recoverer, backOff);
    errorHandler.addNotRetryableExceptions(
            DeserializationException.class, IllegalArgumentException.class);
    container.setCommonErrorHandler(errorHandler);
    return container;
}

DefaultErrorHandler는 같은 스레드에서 블로킹 재시도를 한다. 재시도 중 해당 파티션의 다른 메시지 처리가 지연되지만 메시지 순서는 보장된다. 흐름은 다음과 같다. 메시지 처리 실패 후 ErrorHandler 진입, NotRetryable 예외면 즉시 Recoverer로, 그 외에는 BackOff 대기 후 재시도, 재시도 소진 시 Recoverer가 topic.DLT로 발행한다.

@RetryableTopic (논블로킹 재시도)

@RetryableTopic(
        attempts = "4",
        backoff = @Backoff(delay = 1000, multiplier = 2.0, maxDelay = 10000),
        topicSuffixingStrategy = TopicSuffixingStrategy.SUFFIX_WITH_INDEX_VALUE,
        dltTopicSuffix = "-dlt",
        include = RuntimeException.class,
        exclude = {DeserializationException.class, IllegalArgumentException.class}
)
@KafkaListener(topics = "order-topic", groupId = "order-consumer-group")
public void listen(ConsumerRecord<String, OrderEvent> record) {
    orderService.process(record.value());
}
 
@DltHandler
public void handleDlt(ConsumerRecord<String, OrderEvent> record) {
    notificationService.notifyFailure(record);
}

Spring Kafka가 order-topic-retry-0/1/2, order-topic-dlt 토픽을 자동 생성한다. 실패한 메시지는 재시도 토픽으로 발행되고, 원본 토픽의 다른 메시지는 계속 처리된다 (논블로킹).

ErrorHandler vs @RetryableTopic 선택

측면DefaultErrorHandler@RetryableTopic
방식동기 (블로킹)비동기 (별도 토픽)
순서 보장보장됨깨질 수 있음
처리량재시도 중 정체정상 메시지 계속 처리
토픽 수원본 + DLT원본 + 재시도N + DLT
복잡도낮음중간
적합 케이스순서 중요처리량 중요

처리량 증대 전략

단일 인스턴스의 처리량을 늘리는 접근은 세 가지다.

  1. 배치 컨슈머: max.poll.recordsfetch.min.bytes를 크게 두고 setBatchListener(true)로 한 번에 많이 가져온다. 가장 효율이 좋다.
@Bean
fun batchListenerContainerFactory(): ConcurrentKafkaListenerContainerFactory<String, Any> =
    ConcurrentKafkaListenerContainerFactory<String, Any>().apply {
        consumerFactory = DefaultKafkaConsumerFactory(mapOf(
            ConsumerConfig.MAX_POLL_RECORDS_CONFIG to 500,
            ConsumerConfig.FETCH_MIN_BYTES_CONFIG to 1_048_576,
            ConsumerConfig.FETCH_MAX_WAIT_MS_CONFIG to 500,
        ))
        setConcurrency(10)
        isBatchListener = true
    }
 
@KafkaListener(topics = ["events"], containerFactory = "batchListenerContainerFactory")
fun listen(records: List<ConsumerRecord<String, String>>, ack: Acknowledgment) {
    process(records.map { it.value() })
    ack.acknowledge()
}
  1. 비동기 컨슈머 (ThreadPoolExecutor): 단일 파티션에서 읽은 메시지를 별도 워커 스레드로 분배한다.
@KafkaListener(topics = ["events"])
fun listen(records: List<ConsumerRecord<String, Event>>, ack: Acknowledgment) {
    val futures = records.map { it.value() }.chunked(50).map { chunk ->
        CompletableFuture.runAsync({ eventProcessor.processBatch(chunk) }, executor)
    }
    CompletableFuture.allOf(*futures.toTypedArray()).join()
    ack.acknowledge()
}

장점은 단일 파티션 처리량 증가. 단점은 순서 보장 불가와 부분 실패 시 복구 복잡도.

  1. Confluent Parallel Consumer: 단일 파티션 단일 컨슈머가 내부적으로 여러 워커 스레드를 두고 병렬 처리한다. 미처리 오프셋을 추적해 안전한 지점까지만 커밋한다.
ParallelConsumerOptions.builder<String, String>()
    .consumer(consumer)
    .ordering(ProcessingOrder.KEY)   // PARTITION / KEY / UNORDERED
    .maxConcurrency(100)
    .build()

처리량 비교 (100만 건 기준)

컨슈머 전략처리 옵션TPS소요 시간순위
일반 컨슈머concurrency 10, partition 1018,86753s3
배치 컨슈머concurrency 10, partition 10, batch 50050,00020s1
비동기 컨슈머concurrency 1, partition 1, batch 500, threads 1032,25831s2
병렬 컨슈머partition 1, max concurrency 10, KEY 순서15,62564s4

배치 컨슈머가 가장 빠르다. 비동기 컨슈머는 파티션 1개만 써도 일반 컨슈머보다 낫다. 병렬 컨슈머는 단일 파티션 제약 하에서 일반 컨슈머와 비슷하다.

멀티스레드 처리 패턴

  • 패턴 1 (컨슈머 per 스레드): concurrency를 늘려 컨슈머 인스턴스마다 별도 스레드. 구현 단순, 파티션 단위 순서 보장. 파티션 수 = 최대 병렬도.
  • 패턴 2 (단일 컨슈머 + 워커 풀): 단일 컨슈머가 poll한 결과를 워커 스레드에 분배. 파티션 수보다 많은 병렬 처리 가능하지만 오프셋 커밋과 순서 보장이 복잡해진다.
  • 패턴 3 (pause/resume 배압): 처리가 밀리면 consumer.pause(partitions)로 fetch 중단, 워커 큐가 비면 consumer.resume(). max.poll.interval.ms 초과 방지에 유용하다.

용도별 권장 설정

실시간 처리 (낮은 지연)

spring:
  kafka:
    consumer:
      auto-offset-reset: latest
      enable-auto-commit: false
      properties:
        fetch.min.bytes: 1
        fetch.max.wait.ms: 100
        max.poll.records: 100

배치 처리 (높은 처리량)

spring:
  kafka:
    consumer:
      auto-offset-reset: earliest
      enable-auto-commit: false
      properties:
        fetch.min.bytes: 65536
        fetch.max.wait.ms: 500
        max.poll.records: 1000
        max.partition.fetch.bytes: 1048576

안정적 운영 (공통)

spring:
  kafka:
    consumer:
      enable-auto-commit: false
      properties:
        partition.assignment.strategy: org.apache.kafka.clients.consumer.CooperativeStickyAssignor
        session.timeout.ms: 45000
        heartbeat.interval.ms: 15000
        # max.poll.interval.ms = 처리 시간의 2~3배