카프카의 심장, 토픽·파티션·프로듀서·컨슈머 완벽 해부
1.1 들어가며: 카프카 작동 원리의 핵심#
카프카가 어떻게 대용량 데이터를 실시간으로 처리하고 안정적으로 전달하는지 이해하려면 카프카를 구성하는 핵심 요소부터 알아야 합니다.
카프카는 여러 구성 요소가 서로 협력하면서 데이터를 저장하고 전달하는 분산 이벤트 스트리밍 플랫폼입니다. 그중에서도 브로커, 토픽, 파티션, 프로듀서, 컨슈머는 카프카의 기본 동작을 이해하기 위한 핵심 요소입니다.
이번 장에서는 다음과 같은 내용을 살펴봅니다.
- 브로커: 데이터를 저장하고 관리하는 서버
- 토픽: 데이터를 논리적으로 분류하는 단위
- 파티션: 토픽을 분할하여 병렬 처리와 확장성을 제공하는 단위
- 프로듀서: 카프카로 데이터를 전송하는 애플리케이션
- 컨슈머: 카프카의 데이터를 읽어 사용하는 애플리케이션
- 컨슈머 그룹: 여러 컨슈머가 데이터를 분산 처리하는 방식
- 오프셋: 컨슈머의 데이터 처리 위치를 관리하는 기준
- 멱등성: 프로듀서의 중복 전송을 방지하는 기능
- 복제: 브로커 장애에 대비하여 파티션 데이터를 복제하는 방식
이러한 개념을 이해하면 카프카의 전체 데이터 흐름과 분산 처리 구조를 보다 쉽게 이해할 수 있습니다.
1.2 카프카의 데이터 흐름#
카프카에서 데이터는 프로듀서가 생성하고 카프카의 토픽에 기록됩니다. 토픽은 하나 이상의 파티션으로 구성되며, 실제 메시지는 파티션에 저장됩니다.
브로커는 이러한 파티션을 저장하고 관리하며, 컨슈머는 브로커에 저장된 파티션의 데이터를 읽습니다.
전체적인 흐름은 다음과 같습니다.
[프로듀서]
│
▼
[토픽]
│
▼
[파티션 1] [파티션 2] [파티션 3]
│ │ │
└──────────┼──────────┘
▼
[카프카 브로커]
│
▼
[컨슈머 그룹]
│ │ │
▼ ▼ ▼
[컨슈머] [컨슈머] [컨슈머]중요한 점은 토픽 자체에 메시지가 직접 저장되는 것이 아니라 토픽에 속한 파티션에 메시지가 저장된다는 것입니다.
또한 파티션은 여러 브로커에 분산될 수 있기 때문에 카프카는 많은 양의 데이터를 여러 서버에서 병렬로 처리할 수 있습니다.
1.3 카프카의 핵심 요소#
카프카의 기본 구조를 이해하기 위해서는 브로커, 토픽, 파티션, 프로듀서, 컨슈머의 관계를 먼저 이해해야 합니다.
1.3.1 브로커#
브로커는 카프카 클러스터를 구성하는 서버입니다.
브로커는 프로듀서가 전송한 메시지를 파티션에 저장하고, 컨슈머의 요청에 따라 데이터를 전달합니다.
주요 역할은 다음과 같습니다.
- 토픽 파티션 저장
- 메시지 저장 및 관리
- 프로듀서 요청 처리
- 컨슈머 요청 처리
- 파티션 복제 관리
- 클러스터 구성 관리
카프카는 여러 브로커로 클러스터를 구성할 수 있습니다. 데이터를 여러 브로커에 분산하면 처리량과 저장 용량을 확장할 수 있습니다.
1.3.2 토픽#
토픽은 카프카에서 데이터를 논리적으로 분류하는 단위입니다.
예를 들어 전자상거래 시스템에서는 다음과 같이 토픽을 구성할 수 있습니다.
order
payment
delivery
user
product각 토픽은 특정 업무 영역의 이벤트를 모아서 관리합니다.
토픽은 하나 이상의 파티션으로 구성되며, 실제 데이터는 파티션에 저장됩니다.
1.3.3 파티션#
파티션은 토픽을 여러 개의 독립적인 로그로 나누는 단위입니다.
예를 들어 order 토픽을 3개의 파티션으로 구성할 수 있습니다.
order
├── partition-0
├── partition-1
└── partition-2파티션을 사용하면 데이터를 여러 곳으로 분산하고 동시에 처리할 수 있습니다.
파티션의 주요 특징은 다음과 같습니다.
- 병렬 처리 지원
- 데이터 분산
- 처리량 확장
- 파티션 내부의 메시지 순서 보장
- 브로커 간 데이터 분산 가능
중요한 점은 메시지 순서는 전체 토픽이 아니라 개별 파티션 내부에서 보장된다는 것입니다.
1.3.4 프로듀서#
프로듀서는 카프카에 메시지를 전송하는 애플리케이션입니다.
예를 들어 웹 서비스에서 새로운 주문이 발생하면 프로듀서가 주문 이벤트를 카프카의 order 토픽으로 전송할 수 있습니다.
프로듀서는 메시지의 키와 값을 포함하여 데이터를 전송할 수 있으며, 메시지가 어느 파티션으로 전달될지 결정하는 과정에도 관여합니다.
주요 역할은 다음과 같습니다.
- 데이터 생성
- 메시지 전송
- 토픽 지정
- 메시지 키 지정
- 직렬화
- 압축
- 전송 결과 확인
1.3.5 컨슈머#
컨슈머는 카프카에 저장된 메시지를 읽어 사용하는 애플리케이션입니다.
예를 들어 order 토픽에서 주문 이벤트를 읽어 결제 시스템이나 배송 시스템에서 사용할 수 있습니다.
주요 역할은 다음과 같습니다.
- 토픽 구독
- 메시지 조회
- 메시지 처리
- 오프셋 관리
- 컨슈머 그룹을 통한 병렬 처리
1.4 프로듀서와 컨슈머의 동작#
프로듀서와 컨슈머는 카프카에서 데이터를 생산하고 소비하는 핵심 애플리케이션입니다.
1.4.1 프로듀서 주요 설정#
프로듀서에서는 다음과 같은 설정을 자주 사용합니다.
bootstrap.servers: 카프카 브로커 접속 주소key.serializer: 메시지 키 직렬화 방식value.serializer: 메시지 값 직렬화 방식acks: 메시지 저장 확인 수준enable.idempotence: 프로듀서 멱등성 활성화 여부
1.4.2 자바 프로듀서 예제#
import org.apache.kafka.clients.producer.*;
import org.apache.kafka.common.serialization.StringSerializer;
import java.util.Properties;
public class KafkaProducerExample {
public static void main(String[] args) {
String topicName = "my-topic";
Properties props = new Properties();
props.put(
ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,
"localhost:9092"
);
props.put(
ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,
StringSerializer.class.getName()
);
props.put(
ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,
StringSerializer.class.getName()
);
Producer<String, String> producer =
new KafkaProducer<>(props);
for (int i = 0; i < 10; i++) {
String key = "key-" + i;
String value = "value-" + i;
ProducerRecord<String, String> record =
new ProducerRecord<>(
topicName,
key,
value
);
producer.send(record, (metadata, exception) -> {
if (exception == null) {
System.out.println(
"Message sent to topic "
+ metadata.topic()
+ ", partition "
+ metadata.partition()
+ ", offset "
+ metadata.offset()
);
} else {
System.err.println(
"Failed to send message: "
+ exception.getMessage()
);
}
});
}
producer.flush();
producer.close();
}
}위 코드는 my-topic 토픽에 10개의 메시지를 전송합니다.
각 메시지는 키와 값을 가지며 StringSerializer를 사용하여 문자열 형태로 직렬화됩니다.
producer.send()를 호출하면 메시지가 카프카로 전송되고 콜백을 통해 전송 결과를 확인할 수 있습니다.
1.5 컨슈머의 동작#
컨슈머는 카프카 토픽에 저장된 메시지를 읽습니다.
컨슈머는 일반적으로 컨슈머 그룹에 소속되어 파티션을 나누어 처리합니다.
1.5.1 컨슈머 주요 설정#
bootstrap.servers: 카프카 브로커 접속 주소group.id: 컨슈머 그룹 식별자key.deserializer: 메시지 키 역직렬화 방식value.deserializer: 메시지 값 역직렬화 방식auto.offset.reset: 초기 오프셋 처리 방식
auto.offset.reset은 일반적으로 다음 값을 사용할 수 있습니다.
earliest: 저장된 가장 오래된 메시지부터 읽음latest: 가장 최근 위치부터 읽음none: 기존 오프셋이 없으면 오류 발생
1.5.2 자바 컨슈머 예제#
import org.apache.kafka.clients.consumer.*;
import org.apache.kafka.common.serialization.StringDeserializer;
import java.time.Duration;
import java.util.Collections;
import java.util.Properties;
public class KafkaConsumerExample {
public static void main(String[] args) {
String topicName = "my-topic";
String groupId = "my-group";
Properties props = new Properties();
props.put(
ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG,
"localhost:9092"
);
props.put(
ConsumerConfig.GROUP_ID_CONFIG,
groupId
);
props.put(
ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,
StringDeserializer.class.getName()
);
props.put(
ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,
StringDeserializer.class.getName()
);
props.put(
ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,
"earliest"
);
Consumer<String, String> consumer =
new KafkaConsumer<>(props);
consumer.subscribe(
Collections.singletonList(topicName)
);
while (true) {
ConsumerRecords<String, String> records =
consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> record : records) {
System.out.println(
"Received message: key = "
+ record.key()
+ ", value = "
+ record.value()
+ ", partition = "
+ record.partition()
+ ", offset = "
+ record.offset()
);
}
}
}
}위 코드는 my-topic 토픽을 구독하고 메시지를 계속 읽습니다.
poll() 메서드를 호출하면 컨슈머가 카프카에서 데이터를 가져오며, 각 메시지의 키, 값, 파티션, 오프셋 등을 확인할 수 있습니다.
1.6 컨슈머 그룹과 파티션 할당#
컨슈머 그룹은 여러 컨슈머가 하나의 작업을 나누어 처리하기 위한 논리적인 그룹입니다.
예를 들어 하나의 토픽에 3개의 파티션이 있고 컨슈머 그룹에 3개의 컨슈머가 있다면 다음과 같이 처리할 수 있습니다.
order 토픽
partition-0 ──────> consumer-1
partition-1 ──────> consumer-2
partition-2 ──────> consumer-3이를 통해 여러 컨슈머가 데이터를 병렬로 처리할 수 있습니다.
다만 하나의 컨슈머 그룹에서는 일반적으로 하나의 파티션을 동시에 하나의 컨슈머만 소비합니다.
따라서 파티션 수보다 컨슈머 수가 많으면 일부 컨슈머는 할당받을 파티션이 없어 대기하게 됩니다.
1.6.1 파티션 할당 전략#
카프카는 컨슈머 그룹의 컨슈머에게 파티션을 할당하는 다양한 전략을 제공합니다.
대표적인 전략은 다음과 같습니다.
- Range: 파티션을 범위 단위로 나누어 할당
- RoundRobin: 파티션을 순차적으로 번갈아 할당
- Sticky: 기존 할당을 최대한 유지하면서 재할당
- CooperativeSticky: 기존 할당을 최대한 유지하면서 점진적으로 재할당
파티션과 컨슈머의 수, 토픽 구성에 따라 할당 결과가 달라질 수 있으므로 실제 시스템에서는 환경에 맞는 전략을 선택해야 합니다.
1.7 멱등성#
멱등성은 동일한 작업을 여러 번 수행하더라도 결과가 중복되지 않도록 하는 특성입니다.
카프카에서는 멱등성 프로듀서를 사용하여 네트워크 오류 등으로 인해 프로듀서가 동일한 레코드를 재전송하는 상황에서 브로커가 중복 기록을 방지할 수 있습니다.
프로듀서에서 다음 설정을 사용할 수 있습니다.
enable.idempotence=true멱등성은 컨슈머에게 메시지가 무조건 한 번만 전달되는 기능을 의미하는 것은 아닙니다.
멱등성 프로듀서는 프로듀서와 브로커 사이의 중복 기록을 방지하는 기능이며, 전체 애플리케이션에서 중복 처리를 완전히 없애려면 컨슈머 처리 방식과 트랜잭션 등을 함께 고려해야 합니다.
1.8 오프셋 관리#
오프셋은 파티션에서 메시지를 식별하는 위치 정보입니다.
컨슈머는 오프셋을 기준으로 어디까지 메시지를 처리했는지 관리합니다.
예를 들어 다음과 같은 메시지가 있다고 가정해 보겠습니다.
partition-0
offset 0
offset 1
offset 2
offset 3
offset 4컨슈머가 offset 2까지 처리하고 커밋했다면 이후에는 offset 3부터 이어서 처리할 수 있습니다.
오프셋 관리는 데이터 중복 처리와 데이터 누락에 직접적인 영향을 줄 수 있기 때문에 중요합니다.
1.8.1 자동 커밋#
자동 커밋은 컨슈머가 일정한 주기로 오프셋을 자동으로 커밋하는 방식입니다.
장점:
- 설정이 간단함
- 별도의 커밋 코드가 적음
주의점:
- 메시지 처리 완료 시점과 오프셋 커밋 시점이 일치하지 않을 수 있음
- 애플리케이션 장애 상황에서 메시지 재처리 또는 누락 가능성을 고려해야 함
1.8.2 수동 커밋#
수동 커밋은 애플리케이션이 적절한 시점에 오프셋을 직접 커밋하는 방식입니다.
일반적으로 메시지 처리 완료 후 오프셋을 커밋하도록 구성하여 처리 흐름을 보다 세밀하게 제어할 수 있습니다.
다만 처리와 커밋 사이에 장애가 발생하면 동일한 메시지가 다시 처리될 수 있으므로 애플리케이션의 중복 처리 가능성도 함께 고려해야 합니다.
1.8.3 정확히 한 번 처리#
정확히 한 번 처리는 메시지가 처리 과정에서 중복 효과를 발생시키지 않도록 보장하는 처리 모델입니다.
카프카에서는 멱등성 프로듀서와 트랜잭션 기능 등을 조합하여 정확히 한 번 처리 의미를 구현할 수 있습니다.
1.9 데이터 복제와 고가용성#
카프카는 브로커 장애에 대응하기 위해 파티션 데이터를 여러 브로커에 복제할 수 있습니다.
이를 위해 사용하는 핵심 개념이 복제 팩터입니다.
예를 들어 복제 팩터가 3이면 하나의 파티션에 대한 복제본을 3개 유지할 수 있습니다.
Broker 1
└── Partition 0
└── Leader
Broker 2
└── Partition 0
└── Follower
Broker 3
└── Partition 0
└── Follower1.9.1 리더#
리더는 해당 파티션의 읽기와 쓰기를 처리하는 역할을 담당합니다.
프로듀서와 컨슈머는 파티션의 리더 정보를 바탕으로 데이터를 처리합니다.
1.9.2 팔로워#
팔로워는 리더의 데이터를 복제하여 보관합니다.
리더에 장애가 발생하면 조건을 만족하는 팔로워가 새로운 리더로 선출될 수 있습니다.
이를 통해 특정 브로커에 장애가 발생하더라도 서비스를 계속 운영할 수 있습니다.
1.10 카프카의 확장성과 고가용성#
카프카는 파티션과 브로커를 활용하여 대규모 데이터를 분산 처리할 수 있습니다.
1.10.1 확장성#
데이터가 증가하면 파티션과 브로커를 적절하게 구성하여 처리 능력을 확장할 수 있습니다.
Kafka Cluster
┌───────────────┐
│ Broker 1 │
│ Partition 0 │
└───────────────┘
┌───────────────┐
│ Broker 2 │
│ Partition 1 │
└───────────────┘
┌───────────────┐
│ Broker 3 │
│ Partition 2 │
└───────────────┘파티션을 여러 브로커에 분산하면 데이터를 병렬로 처리할 수 있으며, 컨슈머 그룹을 활용하면 여러 컨슈머가 파티션을 나누어 처리할 수 있습니다.
1.10.2 고가용성#
고가용성은 브로커 장애가 발생하더라도 서비스를 계속 제공할 수 있도록 구성하는 것을 의미합니다.
카프카에서는 파티션 복제와 리더 선출 등을 활용하여 장애 상황에 대응합니다.
따라서 실제 운영 환경에서는 다음 요소를 함께 고려해야 합니다.
- 브로커 수
- 파티션 수
- 복제 팩터
- 장애 허용 수준
- 데이터 보존 기간
- 컨슈머 처리량
- 디스크 용량
- 네트워크 처리량
1.11 전체 데이터 흐름 정리#
카프카의 전체적인 데이터 흐름은 다음과 같이 이해할 수 있습니다.
┌────────────┐
│ Producer │
└─────┬──────┘
│ 메시지 전송
▼
┌────────────┐
│ Topic │
└─────┬──────┘
│
▼
┌─────────────────────────┐
│ Partitions │
│ │
│ Partition 0 │
│ Partition 1 │
│ Partition 2 │
└───────────┬─────────────┘
│
▼
┌─────────────────────────┐
│ Brokers │
│ │
│ 저장 · 복제 · 관리 │
└───────────┬─────────────┘
│
▼
┌─────────────────────────┐
│ Consumer Group │
│ │
│ Consumer 1 │
│ Consumer 2 │
│ Consumer 3 │
└─────────────────────────┘핵심 관계를 정리하면 다음과 같습니다.
프로듀서
↓
토픽
↓
파티션
↓
브로커
↓
컨슈머 그룹
↓
컨슈머여기서 토픽은 데이터를 논리적으로 분류하고, 파티션은 데이터를 분산하여 저장하고 처리합니다. 브로커는 파티션 데이터를 저장하고 관리하며, 컨슈머 그룹은 파티션을 나누어 데이터를 처리합니다.
1.12 요약#
이번 장에서는 카프카의 핵심 요소와 데이터 흐름을 살펴보았습니다.
핵심 내용을 정리하면 다음과 같습니다.
- 브로커는 카프카 클러스터에서 데이터를 저장하고 관리합니다.
- 토픽은 데이터를 논리적으로 분류하는 단위입니다.
- 파티션은 토픽을 분할하여 병렬 처리와 확장성을 제공합니다.
- 프로듀서는 카프카에 메시지를 전송합니다.
- 컨슈머는 카프카의 메시지를 읽어 처리합니다.
- 컨슈머 그룹은 여러 컨슈머가 파티션을 나누어 처리하도록 합니다.
- 오프셋은 컨슈머의 메시지 처리 위치를 관리하는 기준입니다.
- 멱등성은 프로듀서 재전송에 따른 중복 기록을 방지하는 데 사용됩니다.
- 복제는 브로커 장애에 대응하기 위한 핵심 기능입니다.
- 파티션과 브로커를 적절하게 구성하면 대규모 데이터 처리와 확장이 가능합니다.
카프카를 제대로 활용하려면 각각의 개념을 개별적으로 이해하는 것뿐만 아니라 프로듀서 → 토픽 → 파티션 → 브로커 → 컨슈머 그룹 → 컨슈머로 이어지는 전체 데이터 흐름을 이해하는 것이 중요합니다.
다음 장에서는 카프카를 직접 설치하지 않고도 온라인 환경에서 카프카의 기본 동작을 실습하는 방법을 알아봅니다.
1.13 자체 점검 문제#
- 카프카의 핵심 구성 요소인 브로커, 토픽, 파티션, 프로듀서, 컨슈머의 역할을 설명하십시오.
- 카프카에서 프로듀서가 전송한 데이터가 컨슈머에게 전달되는 과정을 설명하십시오.
- 토픽과 파티션의 관계를 설명하고 파티션을 사용하는 이유를 설명하십시오.
- 컨슈머 그룹의 개념을 설명하고 컨슈머 그룹을 사용하는 이유를 설명하십시오.
- 파티션 할당 전략의 종류와 각각의 특징을 설명하십시오.
- 멱등성의 개념을 설명하고 카프카에서 멱등성 프로듀서가 어떤 역할을 하는지 설명하십시오.
- 오프셋의 개념과 오프셋 관리가 중요한 이유를 설명하십시오.
- 자동 커밋과 수동 커밋의 차이점을 설명하십시오.
- 데이터 복제의 개념을 설명하고 리더와 팔로워의 역할을 설명하십시오.
- 복제 팩터가 고가용성과 데이터 안정성에 어떤 영향을 미치는지 설명하십시오.
- 카프카의 파티션과 브로커가 확장성에 어떤 역할을 하는지 설명하십시오.
- 프로듀서와 컨슈머 코드를 작성하고 실행하여 메시지를 전송하고 소비하는 과정을 설명하십시오.