프로듀서와 컨슈머: 카프카 데이터 송수신의 A to Z
1장: 프로듀서와 컨슈머: 카프카 데이터 송수신의 A to Z#
1.1 들어가며: 카프카 데이터 흐름의 시작과 끝#
카프카를 처음 접하면 토픽과 파티션의 개념부터 배우게 됩니다. 하지만 실제 애플리케이션에서 카프카를 사용하려면 결국 다음과 같은 질문에 도달합니다.
데이터를 카프카에 어떻게 넣을까? 넣은 데이터는 어떻게 가져올까? 여러 서버가 동시에 가져가면 어떻게 될까? 장애가 발생하면 어디까지 다시 처리해야 할까?
이 질문의 중심에 있는 것이 바로 프로듀서와 컨슈머입니다.
프로듀서는 애플리케이션에서 발생한 데이터를 카프카로 전달합니다.
애플리케이션
│
▼
Producer
│
▼
Topic
│
┌───┼───┐
▼ ▼ ▼
P0 P1 P2
│ │ │
└───┼───┘
▼
Consumer
│
▼
애플리케이션프로듀서는 메시지를 특정 토픽으로 전송하고, 카프카는 메시지를 파티션에 저장합니다. 컨슈머는 파티션에 저장된 메시지를 가져와 애플리케이션에서 처리합니다.
즉, 매우 단순하게 표현하면 다음과 같습니다.
Producer → Topic → Partition → Consumer하지만 실제 환경에서는 이 흐름에 직렬화, 파티션 선택, 배치, 압축, 재시도, 멱등성, 오프셋, 컨슈머 그룹, 리밸런싱, 장애 복구 등이 추가됩니다.
따라서 프로듀서와 컨슈머를 제대로 이해하려면 단순히 send()와 poll() 메서드만 알아서는 부족합니다.
이 장에서는 프로듀서와 컨슈머가 실제로 어떻게 동작하는지부터 시작하여 성능과 안정성을 결정하는 주요 설정, 컨슈머 그룹과 오프셋 관리, UI 기반 모니터링까지 단계적으로 살펴봅니다.
1.2 프로듀서: 카프카에 메시지를 보내는 주인공#
프로듀서는 애플리케이션에서 생성한 데이터를 카프카 브로커로 전송하는 클라이언트입니다.
예를 들어 쇼핑몰에서 주문이 발생했다고 생각해 봅시다.
사용자 주문
│
▼
주문 서비스
│
▼
Kafka Producer
│
▼
orders 토픽주문 서비스가 직접 결제 서비스, 배송 서비스, 재고 서비스에 데이터를 각각 전달하는 대신 카프카에 주문 이벤트를 저장할 수 있습니다.
┌→ 결제 서비스
│
주문 서비스 → Kafka → 재고 서비스
│
└→ 배송 서비스이 구조를 사용하면 주문 서비스와 각각의 후속 서비스가 느슨하게 연결됩니다.
프로듀서가 수행하는 주요 작업은 다음과 같습니다.
- 카프카 브로커 연결
- 메시지 생성
- 키와 값 직렬화
- 파티션 선택
- 메시지 배치
- 메시지 압축
- 브로커 전송
- 전송 결과 확인
- 실패 시 재시도
1.3 프로듀서의 메시지는 어떻게 만들어질까?#
카프카 메시지는 기본적으로 키와 값으로 구성됩니다.
Key → 사용자 ID
Value → 주문 정보예를 들어 다음과 같은 데이터를 보낼 수 있습니다.
{
"orderId": 10001,
"userId": 200,
"productId": 501,
"amount": 35000
}키를 지정하면 파티션 선택에 활용할 수 있습니다.
Key = user-200
Value = 주문 정보동일한 키를 사용하는 메시지는 일반적으로 동일한 파티션으로 보내지기 때문에 특정 사용자나 주문 등의 이벤트 순서를 유지해야 하는 경우 유용합니다.
다만 카프카에서 순서가 보장되는 범위는 파티션 내부라는 점이 중요합니다.
Partition 0
A → B → C → D
Partition 1
X → Y → Z전체 토픽에 대해 하나의 전역적인 순서가 보장되는 것은 아닙니다.
1.4 프로듀서의 기본 API#
Java에서는 KafkaProducer를 사용하여 메시지를 전송할 수 있습니다.
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()
);
KafkaProducer<String, String> producer =
new KafkaProducer<>(props);
ProducerRecord<String, String> record =
new ProducerRecord<>(
"orders",
"user-100",
"{\"orderId\":10001}"
);
producer.send(record);
producer.close();가장 중요한 흐름은 다음과 같습니다.
ProducerRecord 생성
↓
KafkaProducer.send()
↓
직렬화
↓
파티션 결정
↓
배치에 저장
↓
브로커 전송
↓
응답 수신1.5 프로듀서의 주요 설정#
프로듀서는 상당히 많은 설정을 제공합니다. 그중 실무에서 특히 중요한 설정을 살펴보겠습니다.
1.5.1 bootstrap.servers#
카프카 클러스터에 처음 연결하기 위한 브로커 주소입니다.
bootstrap.servers=broker1:9092,broker2:9092,broker3:9092여러 브로커를 지정하는 것이 일반적입니다.
이 설정은 모든 브로커 목록을 의미하는 것이 아니라 초기 연결에 사용할 주소 목록이라는 점을 이해해야 합니다.
1.5.2 key.serializer#
메시지 키를 바이트 형태로 변환하는 방법을 지정합니다.
key.serializer=org.apache.kafka.common.serialization.StringSerializer1.5.3 value.serializer#
메시지 값을 바이트 형태로 변환합니다.
value.serializer=org.apache.kafka.common.serialization.StringSerializerJSON 객체를 직접 보내는 것이 아니라 실제 네트워크 전송 과정에서는 바이트 형태로 변환되어야 합니다.
Java 객체
↓
Serializer
↓
byte[]
↓
Kafka1.6 acks: 메시지를 얼마나 안전하게 보낼 것인가?#
acks는 프로듀서가 브로커로부터 어느 수준의 확인 응답을 받을 것인지 결정합니다.
1.6.1 acks=0#
브로커의 응답을 기다리지 않습니다.
Producer
│
├──── Message ────→ Broker
│
└── 응답 기다리지 않음빠르지만 메시지가 실제로 저장되었는지 확인할 수 없습니다.
1.6.2 acks=1#
파티션 리더가 메시지를 기록한 후 응답합니다.
Producer
│
▼
Leader
│
└── ACKacks=0보다 안정적이지만 리더 장애와 복제 상황에 따라 데이터 손실 가능성을 완전히 제거하지는 못합니다.
1.6.3 acks=all#
리더뿐만 아니라 현재 ISR에 포함된 복제본들이 메시지를 기록한 조건을 만족해야 응답합니다.
Producer
│
▼
Leader
┌─┴─┐
▼ ▼
F1 F2흔히 acks=all을 가장 높은 수준의 내구성을 제공하는 설정으로 설명하지만, 실제 안정성은 min.insync.replicas 등의 설정과 함께 이해해야 합니다.
1.7 프로듀서의 멱등성#
실무에서 중요한 프로듀서 기능 중 하나가 멱등성입니다.
네트워크 오류가 발생했다고 생각해 봅시다.
Producer
│
│ Message
▼
Broker
│
│ 저장 성공
X
응답 유실프로듀서는 응답을 받지 못했기 때문에 메시지가 실패했다고 판단하고 다시 전송할 수 있습니다.
그러면 동일한 메시지가 중복 기록될 가능성이 있습니다.
멱등성 프로듀서는 이러한 재시도 과정에서 중복 기록을 줄이기 위한 기능을 제공합니다.
일반적으로 다음과 같이 설정할 수 있습니다.
enable.idempotence=true멱등성은 단순한 재시도와 함께 생각해야 합니다.
전송
↓
응답 유실
↓
재시도
↓
중복 가능성멱등성을 활성화하면 프로듀서가 메시지 전송의 중복 문제를 제어하는 데 도움을 줍니다.
다만 멱등성 프로듀서가 애플리케이션 전체의 중복 처리 문제를 해결해 주는 것은 아닙니다.
예를 들어 데이터베이스에 동일한 주문을 두 번 저장하는 문제까지 자동으로 해결해 주지는 않습니다.
1.8 프로듀서 재시도와 전달 시간#
프로듀서에서 실패가 발생하면 메시지를 다시 전송할 수 있습니다.
대표적으로 다음 설정을 함께 살펴볼 필요가 있습니다.
retries
delivery.timeout.ms
request.timeout.msretries는 실패 시 재시도와 관련된 설정이고, delivery.timeout.ms는 메시지가 성공적으로 전송되거나 최종적으로 실패하기까지 허용되는 전체 시간을 관리하는 데 사용됩니다.
따라서 최신 카프카 환경에서는 단순히 retries 하나만 보고 전송 안정성을 판단하기보다는 전체 전달 시간과 재시도 정책을 함께 살펴보는 것이 좋습니다.
1.9 프로듀서 배치 처리#
프로듀서는 메시지를 하나씩 즉시 네트워크로 보내는 대신 여러 메시지를 묶어서 전송할 수 있습니다.
메시지 A ─┐
메시지 B ─┼→ Batch → Broker
메시지 C ─┘이것이 배치 처리입니다.
배치 처리는 네트워크 요청 횟수를 줄이고 처리량을 높이는 데 도움이 됩니다.
주요 설정은 다음과 같습니다.
batch.size
linger.msbatch.size는 하나의 배치에 담을 수 있는 크기와 관련됩니다.
linger.ms는 더 많은 메시지가 배치에 모일 수 있도록 잠시 기다리는 시간과 관련됩니다.
예를 들어 다음과 같은 상황을 생각할 수 있습니다.
메시지 도착
↓
잠시 대기
↓
메시지 추가
↓
Batch 완성
↓
한 번에 전송처리량이 중요한 시스템에서는 배치와 압축을 함께 활용하는 경우가 많습니다.
1.10 메시지 압축#
프로듀서는 메시지를 압축해서 전송할 수 있습니다.
대표적인 압축 방식은 다음과 같습니다.
compression.type=none
compression.type=gzip
compression.type=snappy
compression.type=lz4
compression.type=zstd압축을 사용하면 네트워크로 전달해야 하는 데이터의 양을 줄일 수 있습니다.
압축 전
████████████████████
압축 후
██████████대신 압축과 해제 과정에서 CPU가 사용됩니다.
따라서 압축은 단순히 압축률이 높은 것이 무조건 좋은 것이 아니라 다음 요소를 함께 고려해야 합니다.
- 네트워크 대역폭
- CPU 사용량
- 처리량
- 지연 시간
- 메시지 크기
최근 환경에서는 zstd가 높은 압축 효율과 성능의 균형을 제공하는 선택지로 자주 사용됩니다.
2장: 컨슈머: 카프카에서 데이터를 가져오는 방법#
2.1 컨슈머는 무엇을 하는가?#
컨슈머는 카프카 토픽에 저장된 메시지를 읽어 애플리케이션에서 처리하는 클라이언트입니다.
Kafka Topic
│
▼
Partition
│
▼
Consumer
│
▼
애플리케이션 처리예를 들어 주문 이벤트를 처리하는 컨슈머라면 다음과 같은 작업을 수행할 수 있습니다.
Kafka
↓
주문 이벤트
↓
Consumer
↓
주문 검증
↓
DB 저장
↓
배송 시스템 전달컨슈머의 핵심 API는 poll()입니다.
2.2 컨슈머 기본 코드#
Java에서는 다음과 같은 형태로 컨슈머를 구성할 수 있습니다.
Properties props = new Properties();
props.put(
ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG,
"localhost:9092"
);
props.put(
ConsumerConfig.GROUP_ID_CONFIG,
"order-service"
);
props.put(
ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,
StringDeserializer.class.getName()
);
props.put(
ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,
StringDeserializer.class.getName()
);
KafkaConsumer<String, String> consumer =
new KafkaConsumer<>(props);
consumer.subscribe(List.of("orders"));
while (true) {
ConsumerRecords<String, String> records =
consumer.poll(Duration.ofMillis(1000));
for (ConsumerRecord<String, String> record : records) {
System.out.println(record.value());
}
}전체 흐름은 다음과 같습니다.
Consumer 생성
↓
Topic 구독
↓
poll()
↓
Record 수신
↓
비즈니스 로직 처리
↓
Offset 관리2.3 group.id: 컨슈머 그룹의 시작#
컨슈머를 이해할 때 가장 중요한 설정 중 하나가 group.id입니다.
group.id=order-service같은 group.id를 사용하는 컨슈머들은 하나의 컨슈머 그룹을 구성합니다.
예를 들어 파티션이 3개라면 다음과 같이 처리할 수 있습니다.
Topic
├─ Partition 0 → Consumer A
├─ Partition 1 → Consumer B
└─ Partition 2 → Consumer C컨슈머를 여러 개 실행하면 하나의 애플리케이션이 여러 파티션을 병렬로 처리할 수 있습니다.
3장: 컨슈머 그룹과 스케일 아웃#
3.1 컨슈머 그룹이 필요한 이유#
대량의 데이터를 하나의 컨슈머가 모두 처리한다고 생각해 봅시다.
P0 ─┐
P1 ─┤
P2 ─┼→ Consumer
P3 ─┤
P4 ─┘데이터가 급격하게 증가하면 하나의 컨슈머만으로는 처리량을 감당하기 어려워질 수 있습니다.
컨슈머를 여러 개 실행하면 다음과 같이 분산할 수 있습니다.
P0 → Consumer A
P1 → Consumer B
P2 → Consumer C
P3 → Consumer D이것이 카프카의 대표적인 수평 확장 방식입니다.
3.2 파티션보다 컨슈머가 많으면?#
파티션이 3개인데 컨슈머가 5개라고 생각해 봅시다.
Partition 0 → Consumer A
Partition 1 → Consumer B
Partition 2 → Consumer C
Consumer D → 대기
Consumer E → 대기하나의 컨슈머 그룹에서는 하나의 파티션을 동시에 여러 컨슈머가 소비할 수 없습니다.
따라서 일반적으로 하나의 컨슈머 그룹에서 실제 병렬 처리 수준은 파티션 수에 의해 제한됩니다.
병렬 처리 상한 ≈ 파티션 수이것이 토픽의 파티션 설계가 중요한 이유입니다.
3.3 서로 다른 컨슈머 그룹은 어떻게 동작할까?#
서로 다른 컨슈머 그룹은 동일한 메시지를 독립적으로 소비할 수 있습니다.
Topic
│
┌───────┴───────┐
▼ ▼
주문 그룹 분석 그룹
│ │
주문 처리 통계 처리예를 들어 하나의 orders 토픽을 다음 서비스가 동시에 사용할 수 있습니다.
orders
│
├→ order-service
├→ payment-service
├→ inventory-service
└→ analytics-service각 서비스가 서로 다른 group.id를 사용하면 독립적으로 메시지를 소비할 수 있습니다.
4장: 오프셋: 컨슈머는 어디까지 읽었을까?#
4.1 오프셋이란?#
카프카의 각 파티션에는 메시지마다 순서 번호가 있습니다.
이를 오프셋이라고 합니다.
Partition 0
Offset
0 → A
1 → B
2 → C
3 → D
4 → E컨슈머가 A, B, C까지 처리했다면 현재 처리 위치를 추적할 수 있어야 합니다.
오프셋은 바로 이러한 위치를 나타내는 핵심 개념입니다.
4.2 커밋된 오프셋#
컨슈머 그룹은 처리 위치를 카프카에 커밋할 수 있습니다.
Consumer
│
│ 처리
▼
Message A
Message B
Message C
│
▼
Offset Commit컨슈머가 장애로 재시작되었을 때 커밋된 오프셋을 기반으로 다시 시작할 수 있습니다.
따라서 오프셋 관리는 장애 복구와 데이터 처리 보장에 매우 중요합니다.
5장: 자동 커밋과 수동 커밋#
5.1 자동 커밋#
다음과 같이 설정하면 컨슈머가 오프셋을 자동으로 커밋할 수 있습니다.
enable.auto.commit=true자동 커밋은 사용하기 간편하지만 메시지를 실제 비즈니스 로직에서 처리한 시점과 오프셋 커밋 시점이 정확히 일치하지 않을 수 있습니다.
예를 들어 다음 상황을 생각해 봅시다.
poll()
↓
오프셋 커밋
↓
비즈니스 로직 처리
↓
애플리케이션 장애이 경우 커밋된 위치보다 실제 처리 위치가 뒤처질 수 있습니다.
따라서 자동 커밋을 사용할 때는 메시지 처리와 커밋 타이밍의 관계를 반드시 이해해야 합니다.
5.2 수동 커밋#
수동 커밋에서는 애플리케이션이 오프셋 커밋 시점을 직접 제어합니다.
consumer.commitSync();또는
consumer.commitAsync();와 같은 방식을 사용할 수 있습니다.
일반적인 개념은 다음과 같습니다.
메시지 가져오기
↓
비즈니스 로직 처리
↓
처리 성공
↓
Offset Commit이 구조를 사용하면 처리가 성공한 뒤 어디까지 읽었다고 기록할 것인지를 애플리케이션에서 명확하게 관리할 수 있습니다.
5.3 중복 처리와 메시지 유실#
오프셋 커밋에서는 중복 처리와 데이터 유실을 함께 이해해야 합니다.
처리 전에 커밋하면 다음과 같은 문제가 발생할 수 있습니다.
Offset Commit
↓
비즈니스 처리
↓
장애반대로 처리 후 커밋한다면 다음과 같은 상황에서 중복 처리가 발생할 수 있습니다.
비즈니스 처리 성공
↓
Offset Commit 실패
↓
재시작
↓
같은 메시지 다시 처리따라서 실제 시스템에서는 단순히 중복이 없어야 한다고 생각하기보다 중복 처리가 발생해도 안전한 구조, 즉 멱등적인 비즈니스 로직을 설계하는 것이 중요합니다.
6장: 컨슈머의 주요 설정#
6.1 auto.offset.reset#
컨슈머 그룹에 커밋된 오프셋이 없거나 현재 오프셋을 사용할 수 없는 경우 어디서부터 읽을지를 결정합니다.
대표적인 값은 다음과 같습니다.
auto.offset.reset=earliest가장 오래된 사용 가능한 메시지부터 읽습니다.
auto.offset.reset=latest현재 이후에 새로 들어오는 메시지부터 읽습니다.
auto.offset.reset=none사용할 수 있는 오프셋이 없으면 예외를 발생시킵니다.
중요한 점은 auto.offset.reset이 단순히 새 컨슈머가 추가되었을 때 적용되는 설정이 아니라는 것입니다.
이미 컨슈머 그룹에 유효한 커밋 오프셋이 있다면 그 오프셋이 우선됩니다.
6.2 max.poll.records#
한 번의 poll()에서 반환할 수 있는 최대 레코드 수를 설정합니다.
max.poll.records=500값이 너무 크면 한 번에 처리해야 하는 데이터가 많아지고 처리 시간이 길어질 수 있습니다.
6.3 max.poll.interval.ms#
컨슈머가 poll() 호출 사이에 허용되는 최대 간격과 관련된 중요한 설정입니다.
비즈니스 로직 처리 시간이 너무 길어 이 시간을 초과하면 컨슈머 그룹에서 해당 컨슈머가 정상적으로 작업을 수행하지 못하는 것으로 판단되어 리밸런싱이 발생할 수 있습니다.
따라서 다음과 같은 상황에서 특히 중요합니다.
poll()
↓
대량 데이터 처리
↓
외부 API 호출
↓
DB 작업
↓
처리 지연
↓
다음 poll() 지연컨슈머의 처리 시간이 긴 시스템이라면 max.poll.records와 max.poll.interval.ms를 함께 살펴봐야 합니다.
6.4 fetch.min.bytes와 fetch.max.wait.ms#
컨슈머가 브로커에서 데이터를 가져오는 방식에도 여러 설정이 있습니다.
fetch.min.bytes
fetch.max.wait.msfetch.min.bytes는 브로커가 컨슈머에게 응답할 때 모으려는 최소 데이터 크기와 관련됩니다.
fetch.max.wait.ms는 충분한 데이터가 모이지 않았을 때 브로커가 얼마나 기다릴 수 있는지를 결정합니다.
이를 조절하면 처리량과 지연 시간 사이의 균형을 조정할 수 있습니다.
7장: 리밸런싱: 컨슈머 그룹의 자리 재배치#
7.1 리밸런싱이란?#
컨슈머 그룹의 구성이나 상태가 변경되면 파티션 할당을 다시 계산해야 할 수 있습니다.
이를 리밸런싱이라고 합니다.
예를 들어 다음과 같은 상태에서
P0 → Consumer A
P1 → Consumer B
P2 → Consumer CConsumer B가 장애로 빠지면 파티션을 다시 배분해야 합니다.
P0 → Consumer A
P1 → Consumer C
P2 → Consumer A 또는 C새로운 컨슈머가 추가되는 경우에도 파티션 할당이 변경될 수 있습니다.
7.2 리밸런싱이 중요한 이유#
리밸런싱은 장애 대응과 확장성 측면에서 중요한 기능이지만, 너무 자주 발생하면 애플리케이션 처리에 영향을 줄 수 있습니다.
Consumer 변경
↓
Rebalance
↓
Partition 재할당
↓
처리 재개따라서 컨슈머 그룹을 운영할 때는 컨슈머 수만 확인하는 것이 아니라 리밸런싱 발생 여부와 컨슈머 상태도 함께 관찰해야 합니다.
8장: 컨슈머 성능을 높이는 방법#
컨슈머 성능을 높이는 방법은 단순히 컨슈머 서버를 추가하는 것만이 아닙니다.
다음 요소를 함께 고려해야 합니다.
Partition 수
Consumer 수
poll 처리량
Batch 크기
DB 처리 속도
외부 API 응답 시간
네트워크
Offset Commit특히 카프카의 처리량은 파티션 구조와 밀접한 관계가 있습니다.
예를 들어 파티션이 8개라면 컨슈머 그룹에서도 최대 8개의 파티션을 병렬로 처리할 수 있습니다.
8 Partitions
↓
최대 8개의 병렬 처리 단위다만 실제 처리량은 컨슈머 애플리케이션의 처리 속도와 외부 시스템의 성능에 따라 달라집니다.
9장: 카프카 UI로 데이터를 직접 확인하기#
카프카를 CLI만으로 관리하면 토픽, 파티션, 오프셋, 컨슈머 그룹의 상태를 한눈에 파악하기 어려울 수 있습니다.
그래서 개발 및 운영 환경에서는 Kafka UI 도구가 매우 유용합니다.
대표적인 도구로 다음을 볼 수 있습니다.
- Kafdrop
- Kafbat UI
- Confluent Control Center
9.1 Kafdrop#
Kafdrop은 가볍게 카프카 클러스터를 확인하기 좋은 웹 UI입니다.
주요 기능은 다음과 같습니다.
Broker
├─ Topic
│ ├─ Partition
│ └─ Message
│
└─ Consumer Group
└─ Lag개발 환경에서 카프카 내부를 빠르게 확인할 때 특히 편리합니다.
9.2 Kafbat UI#
Kafbat UI는 오픈소스 Kafka UI 도구입니다.
여러 카프카 클러스터를 웹에서 관리하고 토픽, 메시지, 컨슈머 그룹 등을 확인할 수 있습니다.
개발 환경에서 카프카를 직접 다루는 경우 Kafdrop과 함께 비교해 볼 만한 도구입니다.
9.3 Confluent Control Center#
Confluent Control Center는 Confluent Platform 환경에서 사용하는 웹 기반 관리 및 모니터링 도구입니다.
클러스터 상태뿐만 아니라 토픽, 메시지, Schema Registry, 컨슈머 그룹 lag 등을 확인할 수 있으며 Kafka Connect와 ksqlDB 관련 기능도 제공합니다.
10장: Consumer Lag을 읽어보자#
카프카 운영에서 매우 중요한 지표가 Consumer Lag입니다.
예를 들어 파티션에 현재까지 다음 위치까지 메시지가 들어왔다고 생각해 봅시다.
Log End Offset = 1000컨슈머가 다음 위치를 처리하고 있다면
Consumer Offset = 850대략적인 차이는 다음과 같습니다.
1000 - 850 = 150이 값이 소비 지연을 판단하는 중요한 지표가 됩니다.
Producer
│
▼
1000 ────────────────→ 최신 데이터
Consumer
│
▼
850 ────────────────→ 현재 처리 위치
← Lag 150 →Lag이 지속적으로 증가한다면 생산 속도가 소비 속도보다 빠르다는 신호일 수 있습니다.
Producer 처리량 > Consumer 처리량이 경우 다음과 같은 원인을 확인해야 합니다.
- 컨슈머 처리 속도 부족
- DB 성능 문제
- 외부 API 지연
- 컨슈머 수 부족
- 파티션 설계 문제
- 네트워크 문제
- 리밸런싱
- 애플리케이션 성능 문제
11장: 역직렬화 오류와 메시지 형식 문제#
프로듀서가 데이터를 카프카에 저장할 때 직렬화가 필요했다면 컨슈머는 반대로 역직렬화를 수행합니다.
Producer
Java Object
↓
Serializer
↓
byte[]
↓
Kafka
↓
byte[]
↓
Deserializer
↓
Java Object
Consumer이 과정에서 문제가 발생하면 컨슈머가 메시지를 정상적으로 처리하지 못할 수 있습니다.
11.1 잘못된 역직렬화 클래스#
예를 들어 프로듀서는 문자열을 보내는데 컨슈머가 다른 형식의 역직렬화기를 사용하는 경우입니다.
Producer
StringSerializer
↓
Kafka
↓
Consumer
잘못된 Deserializer11.2 메시지 형식 불일치#
프로듀서가 JSON을 보냈는데 컨슈머가 예상하는 구조가 다른 경우에도 문제가 발생할 수 있습니다.
{
"userId": 100,
"name": "Kim"
}프로듀서가 필드 구조를 변경했는데 컨슈머가 이전 구조만 알고 있다면 애플리케이션에서 문제가 발생할 수 있습니다.
12장: 스키마 관리가 중요한 이유#
대규모 카프카 시스템에서는 메시지 형식을 안정적으로 관리하는 것이 중요합니다.
특히 다음과 같은 환경에서는 스키마 관리가 중요합니다.
서비스 A
│
▼
Kafka
│
├→ 서비스 B
├→ 서비스 C
├→ 서비스 D
└→ 분석 시스템하나의 메시지 형식을 여러 서비스가 사용하고 있기 때문입니다.
대표적으로 다음과 같은 기술을 사용할 수 있습니다.
- Avro
- Protobuf
- JSON Schema
- Schema Registry
특히 Confluent 생태계에서는 Schema Registry를 이용하여 메시지 스키마와 호환성을 관리할 수 있습니다.
13장: 프로듀서와 컨슈머를 함께 이해하기#
지금까지 살펴본 내용을 하나의 흐름으로 합쳐보겠습니다.
애플리케이션
│
▼
Producer
│
├─ Serializer
│
├─ Partition 선택
│
├─ Batch
│
├─ Compression
│
└─ Retry / Idempotence
│
▼
Kafka
│
┌───┼────┐
▼ ▼ ▼
P0 P1 P2
│ │ │
└───┼────┘
▼
Consumer Group
│
├─ Consumer A
├─ Consumer B
└─ Consumer C
│
▼
poll()
│
▼
Deserializer
│
▼
비즈니스 로직
│
▼
Offset Commit이 구조를 이해하면 카프카의 핵심 동작 원리가 훨씬 명확해집니다.
14장: 실무에서 기억해야 할 핵심 설정#
프로듀서에서는 다음 설정을 우선적으로 이해하는 것이 좋습니다.
bootstrap.servers
key.serializer
value.serializer
acks
enable.idempotence
retries
delivery.timeout.ms
batch.size
linger.ms
compression.type컨슈머에서는 다음 설정이 중요합니다.
bootstrap.servers
group.id
key.deserializer
value.deserializer
auto.offset.reset
enable.auto.commit
max.poll.records
max.poll.interval.ms
fetch.min.bytes
fetch.max.wait.ms모든 설정을 무조건 변경할 필요는 없습니다.
중요한 것은 애플리케이션의 처리 특성에 맞게 설정을 선택하는 것입니다.
15장: 프로듀서와 컨슈머 성능 최적화의 기본 원칙#
카프카 성능을 개선할 때는 다음과 같은 순서로 생각하면 좋습니다.
1. 파티션 구조 확인
↓
2. Producer 처리량 확인
↓
3. Batch / Compression 확인
↓
4. Consumer 처리량 확인
↓
5. Consumer Lag 확인
↓
6. DB / 외부 시스템 확인
↓
7. 필요한 경우 Consumer Scale Out특히 Consumer Lag만 보고 무조건 컨슈머 수를 늘리는 것은 좋은 접근이 아닙니다.
예를 들어 컨슈머가 DB에 데이터를 저장하느라 느린 상황이라면 컨슈머를 추가해도 결국 DB가 병목이 될 수 있습니다.
Kafka
↓
Consumer × 20
↓
DB
↓
병목따라서 카프카 성능 문제는 전체 데이터 파이프라인의 병목 지점을 찾아 해결하는 방식으로 접근해야 합니다.
16장: 요약#
이번 장에서는 카프카에서 데이터를 실제로 주고받는 핵심 구성 요소인 프로듀서와 컨슈머를 살펴보았습니다.
핵심 내용을 정리하면 다음과 같습니다.
- 프로듀서는 애플리케이션의 데이터를 카프카 토픽으로 전송합니다.
- 프로듀서는 직렬화, 파티션 선택, 배치, 압축, 재시도 등의 과정을 거쳐 메시지를 전송합니다.
acks는 브로커의 메시지 수신 확인 수준을 결정합니다.- 멱등성 프로듀서는 재시도 과정에서 발생할 수 있는 중복 기록 문제를 줄이는 데 도움을 줍니다.
batch.size와linger.ms는 프로듀서의 배치 처리에 중요한 설정입니다.- 컨슈머는 카프카 토픽에서 메시지를 가져와 애플리케이션에서 처리합니다.
- 컨슈머는
poll()을 통해 메시지를 가져옵니다. - 컨슈머 그룹을 이용하면 여러 컨슈머가 파티션을 나누어 병렬 처리할 수 있습니다.
- 하나의 컨슈머 그룹에서는 일반적으로 하나의 파티션을 동시에 여러 컨슈머가 처리하지 않습니다.
- 오프셋은 컨슈머의 처리 위치를 관리하는 핵심 개념입니다.
- 자동 커밋과 수동 커밋은 메시지 처리와 오프셋 기록의 관계를 결정합니다.
- 중복 처리를 완전히 제거하기보다는 중복이 발생해도 안전한 멱등적 애플리케이션 설계가 중요합니다.
max.poll.records와max.poll.interval.ms는 컨슈머 처리 성능과 리밸런싱에 영향을 줄 수 있습니다.- Consumer Lag은 컨슈머가 데이터를 얼마나 뒤처져 소비하고 있는지 판단하는 중요한 운영 지표입니다.
- Kafdrop, Kafbat UI, Confluent Control Center 같은 도구를 사용하면 카프카 상태를 시각적으로 확인할 수 있습니다.
- 메시지 형식이 복잡해질수록 Schema Registry와 같은 스키마 관리 체계의 중요성이 커집니다.
결국 카프카를 잘 사용한다는 것은 단순히 메시지를 보내고 받는 코드를 작성하는 것이 아닙니다.
Producer
↓
Partition
↓
Consumer Group
↓
Consumer
↓
Business Logic
↓
Offset
↓
Lag Monitoring이 전체 흐름을 이해하고 처리량, 지연 시간, 데이터 중복, 장애 복구, 확장성을 함께 설계하는 것이 카프카를 제대로 활용하는 핵심입니다.
다음 장에서는 실제 서비스 운영 관점에서 카프카 클러스터, 브로커, 복제, 장애 대응, 모니터링과 운영 관리를 살펴보겠습니다.
17장: 자체 점검 문제#
17.1 프로듀서와 컨슈머 기본 이해#
- 프로듀서와 컨슈머의 역할을 설명하십시오.
- 카프카 메시지에서 키와 값은 어떤 역할을 하는지 설명하십시오.
- 카프카에서 메시지 순서가 보장되는 범위를 설명하십시오.
acks=0,acks=1,acks=all의 차이를 설명하십시오.- 프로듀서의 멱등성이 필요한 이유를 설명하십시오.
17.2 프로듀서 성능과 설정#
batch.size와linger.ms가 프로듀서 성능에 어떤 영향을 미치는지 설명하십시오.compression.type을 사용하는 이유를 설명하십시오.- 프로듀서에서 재시도가 필요한 이유를 설명하십시오.
delivery.timeout.ms가 필요한 이유를 설명하십시오.
17.3 컨슈머와 컨슈머 그룹#
- 컨슈머 그룹이 필요한 이유를 설명하십시오.
- 하나의 컨슈머 그룹에서 컨슈머 수가 파티션 수보다 많을 때 어떤 현상이 발생하는지 설명하십시오.
- 서로 다른 컨슈머 그룹이 동일한 토픽을 소비할 수 있는 이유를 설명하십시오.
- 컨슈머 리밸런싱이 무엇인지 설명하십시오.
17.4 오프셋과 장애 처리#
- 카프카의 오프셋이 무엇인지 설명하십시오.
- 자동 커밋과 수동 커밋의 차이를 설명하십시오.
- 메시지를 처리한 뒤 오프셋 커밋에 실패하면 어떤 문제가 발생할 수 있는지 설명하십시오.
auto.offset.reset=earliest와latest의 차이를 설명하십시오.
17.5 운영과 문제 해결#
max.poll.records가 컨슈머 처리에 어떤 영향을 줄 수 있는지 설명하십시오.max.poll.interval.ms가 컨슈머 그룹 운영에서 중요한 이유를 설명하십시오.- Consumer Lag이 증가하는 원인을 세 가지 이상 설명하십시오.
- 역직렬화 오류가 발생하는 대표적인 원인을 설명하십시오.
- Kafdrop, Kafbat UI, Confluent Control Center를 이용해 카프카를 모니터링할 때 확인할 수 있는 정보를 설명하십시오.
- 프로듀서와 컨슈머 사이에서 메시지 스키마를 관리해야 하는 이유를 설명하십시오.
- Schema Registry를 사용하는 목적을 설명하십시오.