[논문 리뷰] Kafka: a Distributed Messaging System for Log Processing

“Kafka: a Distributed Messaging System for Log Processing” 논문을 개인 공부 및 리뷰를 위해 정리한 글이다.
대규모 분산 환경에서 발생하는 방대한 양의 로그 및 사용자 활동 데이터를 실시간으로 수집하고 처리하는 것은 데이터 인프라의 핵심 과제다. 기존 엔터프라이즈 메시징 시스템은 과도한 전달 보증과 임의 접근 인덱스로 인해 처리량(throughput) 한계에 직면했고, 배치 기반 수집기는 수 초 단위의 실시간 처리를 지원하지 못했다. 링크드인(LinkedIn) 연구진은 이러한 병목을 해결하기 위해 단순화된 저장소 구조, OS 페이지 캐시 및 무복사(zero-copy) 전송, 분산 컨슈머 조율을 결합한 고성능 분산 메시징 시스템 카프카(Kafka)를 제안했다.
- 논문 출처 pdf : Kafka paper
- kafka 공식 문서 : kafka.apache.org/documentation/
1. Introduction
현대 인터넷 서비스에서는 사용자 상호작용과 시스템 상태를 파악하기 위해 매일 방대한 로그 데이터가 발생한다. 이러한 데이터는 단순한 사후 분석용 기록을 넘어 실시간 추천, 검색, 보안 및 모니터링 파이프라인의 핵심 입력으로 활용된다.
로그 데이터는 크게 두 가지 범주로 분류할 수 있다.
- 사용자 활동 이벤트(User activity events): 로그인, 페이지 뷰, 클릭, 좋아요, 공유, 댓글 및 검색 쿼리 등에 해당하는 이벤트
- 운영 및 시스템 메트릭(Operational & System metrics): 서비스 콜 스택, 통화 지연 시간, 오류 같은 운영 메트릭과 각 머신의 CPU, 메모리, 네트워크, 디스크 이용률 같은 시스템 메트릭
- 메트릭(metrics): 키-값 쌍으로 캡처된 단순한 숫자 측정치 1
인터넷 애플리케이션에서 이러한 활동 데이터는 생산 데이터 파이프라인의 핵심 구성 요소로 자리 잡았다. 대표적인 활용처는 다음과 같다.
- 검색 관련성(Search relevance): 실시간 쿼리 및 클릭 이벤트를 검색 랭킹에 반영
- 추천 시스템: 항목 인기나 활동 스트림의 동시 발생 패턴에 기반한 개인화 추천
- 광고 타겟팅 및 리포팅: 광고 노출 대비 클릭 성과 분석 및 과금 보고
- 보안 및 남용 방지: 비정상적인 스크래핑이나 악의적인 남용 행위 탐지
- 소셜 피드 집계: 친구 연결 및 상태 업데이트 액션을 실시간 뉴스피드로 집계
로그 데이터의 급증과 실시간 활용 요구는 기존 데이터 인프라에 심각한 병목을 초래했다. 검색, 추천, 광고 도메인에서는 세분화된 클릭률(CTR) 계산을 위해 사용자가 클릭한 항목뿐만 아니라 노출되었으나 클릭하지 않은 수십 개 항목에 대해서도 레코드를 기록하기 때문에 생성되는 로그의 양이 기하급수적으로 증가한다.
링크드인은 프로덕션 환경에서 실시간 애플리케이션을 수 초 이내의 지연 시간으로 안정적으로 지원해야 함을 확인했다. 이에 전통적인 로그 수집기(log aggregator)의 확장성과 메시징 시스템의 실시간성을 결합한 신규 분산 메시징 시스템인 카프카(Kafka)를 설계 및 구축했다.
카프카는 높은 처리량과 분산 확장성을 기본 설계 목표로 삼고, 온라인과 오프라인 처리를 단일 플랫폼으로 통합했다.
카프카는 발행-구독(pub/sub) 인터페이스를 제공하면서도 파일 시스템 기반의 효율적인 추가 전용 로그 구조를 채택하여, 실시간 컨슈머뿐만 아니라 배치 분석 시스템까지 하나의 인프라로 단순화했다.
2. Related Work
기존 메시징 시스템은 비동기 데이터 흐름을 중계하는 이벤트 버스(event bus) 역할을 수행해 왔으나, 대규모 로그 처리에 적용하기에는 구조적 한계가 존재했다.
기존 엔터프라이즈 메시징 시스템(JMS, IBM Websphere MQ, RabbitMQ 등)이 대규모 로그 처리에 적합하지 않았던 주요 원인은 다음과 같다.
- 과도한 전달 보증(Delivery Guarantee): 기존 시스템은 트랜잭션, 복잡한 필터링, 개별 메시지 승인(acknowledgement) 등 고비용 전달 보증에 집중했다. 그러나 로그 수집에는 이러한 기능이 과잉 설계(overkill)이며, 시스템 기본 구현과 API의 복잡도만 불필요하게 증가시킨다.
- 처리량보다 기능 중심 설계: 기존 시스템은 풍부한 기능을 위해 메시지당 처리 오버헤드를 감수했으나, 로그 스트리밍에서 가장 중요한 제약 조건은 초당 수백만 건을 버텨내는 처리량(throughput)이다.
- 취약한 분산 지원: 여러 머신에 걸쳐 메시지를 파티셔닝하고 분산 저장하는 기능이 미흡하여 수평적 확장(scale-out)이 어려웠다.
- 메시지 축적 시 성능 저하: 큐에 메시지가 쌓이지 않고 즉시 소비되는 환경을 전제로 설계되어, 주기적으로 대량의 데이터를 배치 처리하는 오프라인 컨슈머가 메시지를 축적할 경우 심각한 성능 저하가 발생했다.
또한 메시지 전달 방식에서도 기존 시스템의 푸시(push) 모델 대신 컨슈머 주도형 풀(pull) 모델의 필요성이 대두되었다.
링크드인은 각 컨슈머가 자신의 처리 능력에 맞추어 메시지를 가져오고, 브로커가 처리 한도를 초과하여 데이터를 밀어넣음으로써 발생하는 플러딩(flooding) 현상을 원천 방지하기 위해 pull 기반 아키텍처를 채택했다.
3. Kafka Architecture and Design Principles
카프카는 높은 처리량과 분산 확장을 달성하기 위해 최소한의 단순한 추상화와 직관적인 API 설계를 기반으로 구축되었다.
카프카의 기본 핵심 개념은 다음과 같이 정리된다.
- 토픽(topic): 특정 유형의 메시지 스트림을 구분하는 논리적 단위 2
- 프로듀서(producer): 특정 토픽에 메시지를 발행(publish)하는 주체
- 브로커(broker): 발행된 메시지를 디스크에 저장하고 서빙하는 서버 클러스터 노드
- 컨슈머(consumer): 하나 이상의 토픽을 구독(subscribe)하고, 브로커로부터 데이터를 pull 방식으로 소비하는 주체
카프카는 개념 모델을 단순화한 만큼, 개발자가 다루는 API 역시 직관적으로 설계했다.
프로듀서의 샘플 코드는 다음과 같다.
/* Sample producer code */
producer = new Producer(...);
message = new Message("test message str".getBytes());
set = new MessageSet(message);
producer.send("topic1", set);
컨슈머의 샘플 코드와 동작 방식은 다음과 같다.
- 토픽을 구독하기 위해 컨슈머는 토픽을 위한 하나 이상의 메시지 스트림을 먼저 생성한다.
- 브로커에 발행된 메시지는 여러 서브 스트림으로 균일하게 분산된다.
- 각 메시지 스트림은 생성되는 메시지의 연속적인 스트림에 대해 이터레이터(iterator) 인터페이스를 제공한다.
- 컨슈머는 스트림의 메시지를 순회하며 페이로드(payload)를 꺼내 처리한다.
- 페이로드(payload): 전송 목적이 되는 실제 데이터 본문 3
- 현재 읽을 메시지가 없으면 이터레이터는 새 메시지가 토픽에 발행될 때까지 블로킹된다.
/* Sample consumer code */
streams[] Consumer.createMessageStreams("topic1", 1);
for (message : streams[0]) {
bytes = message.payload();
// do something with the bytes
}
카프카는 여러 컨슈머가 메시지의 단일 복사본을 나누어 처리하는 포인트 투 포인트(point-to-point) 모델과, 각 컨슈머 그룹이 독립적으로 전체 메시지를 소비하는 발행/구독(publish/subscribe) 모델을 모두 지원한다.
- retrieve: 데이터를 검색(searching), 위치 확인(locating), 반환(returning)하는 과정을 의미한다. 4
카프카 클러스터의 전반적인 구조와 파티셔닝 전략은 다음과 같다.
- 카프카 클러스터는 본질적으로 분산 아키텍처이며 여러 대의 브로커로 구성된다.
- 부하 분산과 병렬 처리를 위해 각 토픽을 여러 파티션(partition)으로 분할하며, 브로커마다 하나 이상의 파티션을 분산 저장한다.
- 파티션(partition): 메시지를 순차적으로 기록하는 물리적인 로그 단위 2
- 다수의 프로듀서와 컨슈머가 파티션 단위로 동시에 메시지를 발행하고 검색(retrieve)할 수 있다.

3.1 Efficiency on a Single Partition
카프카가 단일 파티션 수준에서 극단적인 처리량을 달성할 수 있었던 비결은 저장 구조 단순화와 OS 최적화의 적극적 활용에 있다.
Simple storage
토픽의 각 파티션은 하나의 논리적 로그에 대응하며, 물리적으로는 거의 동일한 크기를 갖는 세그먼트(segment) 파일들의 집합으로 디스크에 저장된다.
프로듀서가 메시지를 파티션에 발행할 때마다 브로커는 활성(active) 세그먼트 파일의 끝에 순차적으로 메시지를 추가(append-only)한다. 디스크 쓰기 성능을 극대화하기 위해 설정 가능한(configurable) 메시지 수가 도달하거나 일정 시간이 경과한 후에만 세그먼트 파일을 디스크에 플러시(flush)한다.
- 메시지는 디스크로 플러시된 이후에만 컨슈머에게 노출된다.
- 플러시(flush): 메모리 버퍼의 변경 내용을 영속 저장소에 동기화하여 기록하는 과정 5
일반적인 메시징 시스템과 달리 카프카의 메시지에는 고유 메시지 ID가 부여되지 않으며, 대신 로그 내 논리적 오프셋(offset)으로 식별된다.
임의 접근 B-Tree 대신 순차 로그 오프셋을 사용함으로써 불필요한 인덱스 오버헤드를 완전히 제거했다.
- 오프셋(offset): 파티션 내 각 메시지가 저장된 상대적 위치 번호 2
- 다음 메시지의 위치는 현재 메시지의 길이를 오프셋에 더하여 즉시 계산할 수 있다.
- 이러한 설계 덕분에 메시지 ID와 논리적 오프셋을 동의어로 사용할 수 있다.
컨슈머는 특정 파티션의 메시지를 언제나 순차적으로(sequentially) 소비한다.
컨슈머가 특정 오프셋을 커밋한다는 것은 해당 파티션에서 그 오프셋 이전의 모든 메시지를 정상적으로 수신했음을 의미한다. 컨슈머는 브로커에게 비동기 풀(pull) 요청을 전송하여 로컬 버퍼를 채운다.
- 각 풀 요청에는 읽기를 시작할 메시지 오프셋과 읽어올 최대 바이트 크기가 포함된다.
- 각 브로커는 모든 세그먼트 파일의 시작 오프셋 목록을 메모리상에 정렬된 상태로 유지한다.
- 브로커는 메모리 내 오프셋 목록을 이진 검색하여 요청된 오프셋이 포함된 세그먼트 파일을 빠르게 찾고, 데이터를 컨슈머에게 반환한다.
- 컨슈머는 데이터를 수신한 후 다음 읽을 오프셋을 계산하여 후속 풀 요청에 활용한다.
다음 그림은 카프카의 세그먼트 로그 레이아웃과 메모리 내 오프셋 인덱스의 동작 방식을 보여준다.

Efficient transfer
네트워크 전송과 I/O 효율성을 극대화하기 위해 카프카는 배치 전송, OS 페이지 캐시 활용, 제로 카피(zero-copy) 기법을 도입했다.
최종 컨슈머 API는 메시지를 한 건씩 순회하지만, 네트워크 계층에서는 한 번의 풀(pull) 요청으로 설정된 바이트 한도까지 여러 메시지를 한 번에 묶어서(batch) 가져온다.
또한 카프카는 애플리케이션 레벨의 사용자 공간 메모리 캐싱을 완전히 배제하고, 운영체제의 파일 시스템 페이지 캐시(page cache)에 전적으로 의존한다.
- 이중 버퍼링(Double buffering) 제거: JVM 힙 메모리와 OS 캐시에 동일 데이터가 중복 보관되는 낭비를 없앤다. 메시지는 오직 커널의 페이지 캐시에만 보관된다.
- 웜 캐시(Warm cache) 유지: 브로커 프로세스가 재시작되더라도 OS 페이지 캐시는 그대로 유지되므로 즉각적인 고성능 서빙이 가능하다.
버퍼(buffer) : 어떤 장치에서 다른 장치로 데이터를 송신할 때 일어나는 시간의 차이나 데이터 흐름의 속도 차이를 조정하기 위해 일시적으로 데이터를 기억시키는 장치. 싱글버퍼(single buffer)의 경우 채널이 데이터를 버퍼에 저장하면 프로세서가 처리하는 방식으로 진행된다. 이경우 채널이 데이터를 저장하는 동안에는 데이터에 대한 처리가 이루어질 수 없으며, 프로세서가 데이터를 처리하는 동안에는 다른 데이터가 저장될 수 없게 된다. 이중 버퍼(double buffer)의 경우에는 데이터에 대한 저장과 처리가 동시에 일어날 수 있다. 6
카프카 브로커는 인메모리 객체 캐싱을 하지 않으므로 GC(Garbage Collection) 오버헤드가 거의 없어 JVM 위에서도 안정적인 저지연 처리가 가능하다.
네트워크 전송 시 카프카는 리눅스의 sendfile 시스템 콜을 활용한다.
카프카는 하나의 토픽 데이터를 다수의 컨슈머가 반복 소비하는 멀티 구독자 구조다. 일반적인 소켓 전송은 커널 공간에서 사용자 공간으로 데이터를 복사한 후 다시 소켓 버퍼로 복사하지만, 리눅스의 sendfile API를 활용하면 OS 페이지 캐시에서 네트워크 소켓으로 데이터를 직접 전송하는 제로 카피(zero-copy)가 가능해져 CPU 사용량과 메모리 대역폭 낭비를 획기적으로 줄인다.
Stateless broker
카프카 브로커는 컨슈머의 상태(state)를 일절 기억하지 않는 무상태(stateless) 구조로 설계되어 대규모 동시 접속 상황에서도 오버헤드를 최소화한다.
전통적인 메시징 시스템과 달리, 각 컨슈머가 어디까지 메시지를 읽었는지(오프셋)는 브로커가 아닌 컨슈머가 직접 관리한다. 이러한 설계는 브로커의 복잡도와 상태 관리 오버헤드를 대폭 경감시킨다.
그러나 브로커가 컨슈머의 소비 여부를 추적하지 않으므로, 언제 메시지를 디스크에서 삭제해야 할지가 문제가 된다. 카프카는 이를 시간 기반 SLA(Service Level Agreement) 보존(retention) 정책으로 해결한다.
- 보존 주기 만료 시 자동 삭제: 특정 기간(기본 7일) 이상 보관된 세그먼트 파일은 소비 여부와 무관하게 자동으로 디스크에서 삭제된다.
- 오프라인 컨슈머 지원: 대부분의 실시간 및 일배치 컨슈머는 보존 기간 내에 소비를 완료하므로, 메시지가 쌓여 데이터 크기가 커지더라도 O(1) 디스크 접근 특성 덕분에 카프카의 성능은 저하되지 않는다.
소비 상태를 컨슈머가 제어하므로, 오류 발생 시 이전 오프셋으로 되돌아가(rewind) 데이터를 재처리할 수 있는 강력한 유연성을 제공한다.
3.2 Distributed Coordination
분산 환경에서 수많은 프로듀서, 브로커, 컨슈머가 중앙 집중식 병목 없이 조화롭게 동작하기 위해 주키퍼(ZooKeeper) 기반의 분산 조율 메커니즘을 사용한다.
프로듀서는 메시지 키의 해시 함수를 통하거나 라운드로빈 방식으로 대상 파티션을 결정하여 메시지를 발행한다. 한편 수신 측에서는 컨슈머 그룹(consumer group) 개념을 통해 부하를 분산한다.
- 컨슈머 그룹 분할 소비: 각 컨슈머 그룹은 토픽의 파티션들을 공동으로 나누어 소비한다. 즉, 특정 파티션의 메시지는 그룹 내 단 하나의 컨슈머에게만 전달된다.
- 그룹 간 독립성: 서로 다른 컨슈머 그룹은 독립적으로 전체 토픽 데이터를 온전히 소비하며, 그룹 간에는 어떠한 조율도 요구되지 않는다.
카프카는 복잡한 분산 락(distributed lock)이나 중앙 마스터 노드 없이 파티션을 컨슈머에게 균등하게 분배하기 위해 두 가지 핵심 원칙을 채택했다.
- 파티션을 병렬화의 최소 단위로 설정: 특정 파티션은 동일 컨슈머 그룹 내에서 오직 하나의 컨슈머만 독점하여 소비한다. 이를 통해 메시지 소비 순서를 보장하면서도 컨슈머 간 락 경합을 원천 배제했다.
- 분산 합의 기반의 자율 재조정(Rebalancing): 중앙 집중식 마스터 노드가 파티션을 할당하는 대신, 컨슈머들이 분산 합의 서비스인 주키퍼(ZooKeeper)를 매개로 자율적으로 조율(coordination)하도록 했다.
주키퍼는 트리 구조의 파일 시스템 형태 API를 제공하며, 카프카는 이를 통해 분산 상태를 유지한다.
- 주키퍼(ZooKeeper)의 주요 특성: 경로 생성/조회/삭제, 데이터 변경 시 실시간 통보를 제공하는 워처(watcher), 클라이언트 세션 종료 시 자동으로 삭제되는 임시(ephemeral) znode, 과반수 노드 복제를 통한 고가용성 제공
카프카는 다음과 같은 분산 조율 작업에 주키퍼를 활용한다.
- 브로커와 컨슈머의 동적 추가 및 장애(제거) 탐지
- 멤버십 변경 이벤트 발생 시 각 컨슈머에게 파티션 재조정(rebalance) 프로세스 트리거
- 컨슈머와 파티션 간의 소비 소유권 관계 유지 및 최종 소비 오프셋 추적
주키퍼에 저장되는 레지스트리의 세부 구성은 다음과 같다.
- 브로커 레지스트리(Broker Registry): 브로커의 호스트명, 포트, 해당 브로커가 호스팅하는 토픽 및 파티션 목록을 저장 (임시 노드)
- 컨슈머 레지스트리(Consumer Registry): 컨슈머 그룹 ID와 해당 컨슈머가 구독한 토픽 목록을 저장 (임시 노드)
- 소유권 레지스트리(Ownership Registry): 구독된 파티션마다 현재 메시지를 소비 중인 컨슈머 ID를 매핑 (임시 노드)
- 오프셋 레지스트리(Offset Registry): 컨슈머 그룹별로 각 파티션에서 마지막으로 성공적으로 처리한 커밋 오프셋을 영구 저장 (영구 노드)
브로커나 컨슈머에 장애가 발생하면 주키퍼의 임시 노드가 만료되면서 자동으로 레지스트리에서 제거된다. 각 컨슈머는 브로커 레지스트리와 컨슈머 레지스트리에 워처(watcher)를 등록해 두어, 브로커 구성이나 컨슈머 그룹 멤버가 변경될 때마다 즉시 알림을 수신한다.
컨슈머가 최초 구동되거나 워처 알림을 받으면, 자신이 소비할 파티션 목록을 결정하는 재조정(Rebalance) 알고리즘을 수행한다.

3.3 Delivery Guarantees
로그 데이터의 특성을 고려하여 카프카는 과도한 트랜잭션 오버헤드를 피하고 실용적인 전달 보증 모델을 채택했다.
기본적으로 카프카는 최소 한 번 전달(at-least-once delivery)을 보장한다.
2단계 커밋(2-phase commit)과 같은 고비용의 정확히 한 번(exactly-once) 트랜잭션은 처리량을 저하시키며 로그 분석 시스템에서는 필수적이지 않다.
- 정상적인 상황에서 메시지는 컨슈머 그룹에 정확히 한 번 전달된다.
- 컨슈머 장애 시 파티션 재조정 과정에서 일부 메시지가 중복 처리될 수 있으나, 중복을 허용하지 않는 애플리케이션은 메시지의 고유 키나 오프셋을 활용하여 컨슈머 단에서 자체 멱등성(deduplication)을 구현할 수 있다. 이는 분산 2PC보다 훨씬 가볍고 효율적인 접근법이다.
메시지 순서 보장과 데이터 무결성 검증 정책은 다음과 같이 동작한다.
- 파티션 내 순서 보장: 카프카는 단일 파티션 내부에서 메시지가 발행된 순서대로 컨슈머에게 전달됨을 엄격히 보장한다. 단, 서로 다른 파티션 간의 메시지 순서는 보장하지 않는다.
- CRC 무결성 검증: 디스크 손상이나 전송 중 데이터 오염을 방지하기 위해 각 메시지마다 CRC-32 체크섬을 함께 저장한다.
- 순환 중복 검사(CRC): 전송 또는 저장된 데이터의 오류를 검출하기 위해 계산하는 체크값 7
- I/O 장애 복구: 브로커에 I/O 에러가 발생할 경우 복구 루틴을 실행하여 CRC가 일치하지 않는 손상된 메시지를 로그에서 제거한다.
4. 핵심 요약 및 시사점
카프카는 대규모 분산 로그 처리를 위해 복잡한 기능을 과감히 걷어내고 처리량과 단순성을 극대화한 현대 데이터 엔지니어링의 이정표적 시스템이다.
- 추가 전용 로그(Append-only Log): 복잡한 B-Tree 인덱스 대신 순차 디스크 I/O와 논리적 오프셋을 채택하여 디스크 쓰기 병목을 제거하고 $O(1)$의 일정한 성능을 달성했다.
- OS 커널 기능의 극대화: 자체 인메모리 캐싱을 배제하고 OS 페이지 캐시와
sendfile기반 제로 카피(zero-copy)를 활용해 JVM 가비지 컬렉션 부담을 없애고 네트워크 대역폭을 온전히 활용했다. - 무상태 브로커(Stateless Broker): 소비 오프셋 관리를 컨슈머에게 위임하여 브로커의 상태 오버헤드를 최소화하고, 시간 기반 SLA 보존 정책으로 높은 동시성을 확보했다.
- 파티션 단위 병렬성과 주키퍼 조율: 토픽을 파티션으로 분할하여 병렬성을 극대화하고, 주키퍼를 통한 탈중앙화된 컨슈머 자율 재조정으로 단일 장애점(SPOF) 없는 고가용성 분산 시스템을 완성했다.
댓글남기기