카프카 프로듀서란?
프로듀서는 어플리케이션에서 생성한 레코드를 카프카 브로커로 전송하는 역할을 담당 한다
카프카 프로듀서 흐름

Produer Record
브로커로 보내기 위한 정보를 담은 객체
- 토픽
- 파티션 (option)
- 키 (option)
- 밸류
// 1. 파티션까지 직접 지정하는 경우
ProducerRecord<String, String> record = new ProducerRecord<>("토픽A", 0, "키값", "메시지내용");
// 토픽명 파티션 키 밸류
// 2. 파티션 없이 키만 지정 (파티셔너가 파티션 결정)
ProducerRecord<String, String> record = new ProducerRecord<>("토픽A", "키값", "메시지내용");
// 3. 파티션도 키도 없는 경우
ProducerRecord<String, String> record = new ProducerRecord<>("토픽A", "메시지내용");
카프카 프로듀서 메시지 전송 흐름
- ProducerRecord 생성 (토픽, 파티션, 키, 값 등 포함)
- 직렬화(Serializer): 키와 값을 바이트 배열로 변환
- 파티셔너(Partitioner): 어느 파티션으로 보낼지 결정 (키가 있으면 해시 기반, 없으면 라운드로빈 등)
- 레코드 배치(RecordAccumulator): 같은 파티션으로 가는 레코드를 배치로 묶음
- 전송 스레드가 배치를 브로커로 전송
- 성공 시 RecordMetadata 반환, 실패 시 예외 처리
프로듀서 생성
bootstrap.servers: 브로커 목록- 브로커 동작 이슈 대비하여 최소 2개 이상 권장
key.serializer,value.serializer: 직렬화 클래스 지정- ByteArraySerializer
- StringSerializer
- IntegerSerializer
메시지 전송 방식
파이어 앤 포켓
- 메시지를 전송한 후 결과를 확인 안함
동기적 전송
- 메시지를 전송한 후 반환 되는 Future 객체의 get 메서드를 호출하여 결과를 수신할 때 까지 대기하고 결과 수신 후 다음 처리를 진행
비동기적 전송
- 메시지를 전송 후 브로커로부터 응답을 수신 받을 시점에 callback 함수를 호출하여 별도 비즈니스 로직을 처리
프로듀서 설정
- client.id : 프로듀서와 애플리캐이션을 구분하기 위한 논리적 식별자
acks: 메시지 전달 신뢰성 수준 (0, 1, all/-1)- 0 : 메시지가 성공적으로 전달되었다고 간주하고 브로커의 응답을 기다리지 않음
- 1 : 리더 레플리카가 메시지를 받는 순간 브로커로부터 성공했다는 응답을 받음
- all/-1 : 메시지가 모든 인-싱크 레플리카에 전달된 뒤에 브로커로부터 성공했다는 응답을 받음
buffer.memory: 전송 대기 중인 메시지를 위한 버퍼 크기compression.type: 압축 방식 (snappy, gzip, lz4, zstd)retries/retry.backoff.ms: 재시도 관련 설정batch.size,linger.ms: 배치 크기와 대기 시간 (처리량 vs 지연시간 트레이드오프)max.in.flight.requests.per.connection: 순서 보장과 관련된 설정max.request.size: 최대 요청 크기

delivery.timeout.ms >= linger.ms + retry.backoff.ms + request.timeout.ms
max.block.ms : send를 호출 했을 때, 버퍼가 꽉 차있거나 메타데이터를 아직 못 가져오는 경우 최대 대기 시간
linger.ms : 메시지를 바로 보내지 않고, 배치로 묶기 위해 추가로 기다리는 시간
retry.backoff.ms : 전송이 실패했을 때 재시도하기 전 대기하는 시간
request.timeout.ms : 실제로 브로커에 요청을 보낸 후, 브로커의 응답을 얼마나 기다릴지 정하는 시간
시리얼라이저
기본 시리얼라이저
| 시리얼라이저 | 대상 타입 |
|---|---|
| StringSerializer | String |
| IntegerSerializer | Integer |
| LongSerializer | Long |
| DoubleSerializer | Double |
| ByteArraySerializer | byte[] |
커스텀 시리얼라이저
Serializer인터페이스를 구현하여 Custom 시리얼라이저를 만들 수 있다
public class Customer {
private int customerID;
private String customerName;
// getter, setter 생략
}
public class CustomerSerializer implements Serializer<Customer> {
@Override
public byte[] serialize(String topic, Customer data) {
// customerID(4바이트) + customerName 길이(4바이트) + customerName(바이트)
// 형태로 직접 바이트 배열을 조립
...
}
}
커스텀 시리얼라지어의 문제점
- 스키마 변경에 취약함
- 예:
Customer객체에email필드를 나중에 추가하면, 새 버전으로 보낸 메시지를 예전 버전의 컨슈머(디시리얼라이저)가 읽으면 파싱이 깨지거나 에러 발생 - 프로듀서와 컨슈머 간의 스키마 호환성을 스스로 관리해야 하는 부담이 생김
- 예:
- 유지보수 어려움
- 필드가 추가/삭제될 때마다 시리얼라이저와 디시리얼라이저 코드를 양쪽 다 수정해야 함
- 여러 팀, 여러 서비스가 같은 토픽을 사용하는 경우 동기화가 매우 어려워짐
- 재사용성 부족
- 객체마다 매번 새로운 시리얼라이저 작성 필요
Avro
아파치에서 만든 직렬화 시스템
데이터 구조를 스키마로 별도 정의하고, 스키마를 기반으로 직렬화/역직렬화를 수행
- 스키마와 데이터 분리
- Avro는 데이터를 쓸 때 사용한 스키마와 읽을 때 사용하는 스키마가 완전히 같지 않아도 호환되도록 설계
- 스키마 진화(Schema Evolution) 지원
- 필드 추가/삭제 시에도 하위 호환성(backward compatibility) 또는 상위 호환성(forward compatibility)을 유지 가능
- 스키마 레지스트리(Schema Registry) 연동
- 실무에서는 보통 Avro를 Confluent Schema Registry와 함께 사용
- 동작 방식:
- 프로듀서가 메시지를 보낼 때, 스키마를 레지스트리에 등록 (또는 이미 등록된 스키마 ID 조회)
- 메시지 안에는 스키마 전체가 아니라 스키마 ID만 포함해서 전송 → 메시지 크기 절약
- 컨슈머는 메시지에 포함된 스키마 ID로 레지스트리에서 해당 스키마를 조회해 역직렬화
- 이 구조 덕분에 프로듀서와 컨슈머가 서로 다른 버전의 스키마를 사용해도 자동으로 호환성이 체크 보장
예시
[프로듀서]
Customer 객체 → Avro 스키마 기반 직렬화 → 바이트 + 스키마ID
↓
Schema Registry에 스키마 등록/조회
↓
카프카 브로커로 전송
[컨슈머]
바이트 + 스키마ID 수신
↓
Schema Registry에서 스키마ID로 스키마 조회
↓
Avro 역직렬화 → Customer 객체 복원
파티션
파티셔너(Partitioner)의 역할
프로듀서가 메시지를 보낼 때, 어느 토픽의 몇 번 파티션에 저장할지 결정하는 컴포넌트
시리얼라이저 (바이트 변환)
↓
파티셔너 ← 여기서 파티션 결정
↓
배치 큐 (해당 파티션으로)
- 같은 파티션 내에서는 메시지 순서가 보장
키가 있는 경우 - 해시 기반 파티셔닝
동작 방식
- 기본 파티셔너(DefaultPartitioner)는 키의 해시값을 파티션 개수로 나눈 나머지를 이용해 파티션 결정
partition = hash(key) % 파티션개수
특징
- 같은 키는 항상 같은 파티션으로 감 → 이게 핵심 포인트
- 예: 키가 "user123"인 메시지는 항상 같은 파티션으로 전송됨
- 컨슈머 입장에서 봤을 때 같은 키를 가진 메시지들끼리는 순서가 보장됨
주의할 점
- 파티션 개수가 변경되면 해시 결과가 달라짐 → 기존에 매핑되던 파티션과 달라질 수 있음
- 그래서 토픽의 파티션 수는 운영 중 함부로 늘리지 않는 것이 권장됨 (특히 키 기반 파티셔닝을 사용할 때)
키가 없는 경우 - 라운드로빈 / Sticky 방식
키가 없을 때의 동작
- 키가 없으면 해시 계산을 할 대상이 없으므로, 파티션에 고르게 분산시키는 방식 사용
Sticky 방식으로 바뀌었나
- 순수 라운드로빈은 매번 파티션을 바꾸다 보니, 배치가 작은 상태로 여러 파티션에 흩어져 전송되는 비효율 발생
- 배치 효율 저하 → 처리량 저하
- Sticky Partitioner는 같은 파티션에 메시지를 모아 배치 크기를 키워서 전송 효율(처리량)을 높이는 것이 목적
- 장기적으로 보면 전체 파티션에 고르게 분산되는 건 동일하되, 배치 단위로 묶어서 보내는 효율성 개선
커스텀 파티셔너 (Custom Partitioner)
기본 해시 기반 방식으로는 해결이 안 되는 비즈니스 로직상의 특수한 분배 요구가 있을 때 사용
- Partitioner 인터페이스 구현
헤터, 인터셉터, 쿼터
헤더 (Headers)
카프카 메시지는 키와 밸류 외에 헤더(Header)라는 별도의 키-값 쌍 목록을 가질 수 있다.
메시지 본문(밸류)을 건드리지 않고 부가적인 메타데이터를 함께 전달하고 싶을 때 사용한다.
예시
ProducerRecord<String, String> record =
new ProducerRecord<>("토픽A", "키값", "메시지내용");
record.headers().add("traceId", "abc-123".getBytes(StandardCharsets.UTF_8));
record.headers().add("source", "order-service".getBytes(StandardCharsets.UTF_8));
producer.send(record);
인터셉터 (ProducerInterceptor)
메시지가 전송되기 직전 또는 전송 완료(성공/실패) 직후에 공통 로직을 자동으로 끼워 넣을 수 있는 훅(hook)이다.
애플리케이션 코드를 수정하지 않고도 프로듀서 동작에 공통 처리를 추가할 수 있다.
ProducerInterceptor<K, V>인터페이스 구현- onSend: 시리얼라이저, 파티셔너를 거치기 전 send를 호출한 직후
- onAcknowledgement: 브로커로부터 응답(성공/실패)을 받은 직후
- close: 프로듀서 종료 시
public class CountingProducerInterceptor implements ProducerInterceptor<String, String> {
@Override
public ProducerRecord<String, String> onSend(ProducerRecord<String, String> record) {
// 전송 전: 헤더에 타임스탬프 추가
record.headers().add("sentAt",
String.valueOf(System.currentTimeMillis()).getBytes());
return record;
}
@Override
public void onAcknowledgement(RecordMetadata metadata, Exception exception) {
// 전송 후: 성공/실패 카운트
if (exception != null) {
System.out.println("전송 실패: " + exception.getMessage());
} else {
System.out.println("전송 성공, 오프셋: " + metadata.offset());
}
}
@Override
public void close() {}
@Override
public void configure(Map<String, ?> configs) {}
}
쿼터 (Quotas)
카프카 브로커가 특정 클라이언트(프로듀서/컨슈머)의 처리량을 제한하는 기능
하나의 클라이언트가 브로커 자원을 독점해서 다른 클라이언트에게 영향을 주는 것을 방지하기 위한 장치
쿼터 종류
- 프로듀서 쿼터: 특정 클라이언트가 브로커에 쓸 수 있는 초당 바이트 수 제한
- 컨슈머 쿼터: 특정 클라이언트가 브로커에서 읽을 수 있는 초당 바이트 수 제한
- 요청 쿼터: 브로커 요청 처리 시간(CPU 사용률) 기준 제한
동작 방식
- 클라이언트가 할당된 쿼터를 초과하면, 브로커는 응답을 지연(throttle)시켜서 클라이언트가 자연스럽게 속도를 줄이도록 유도
- 요청을 거부하는 게 아니라 응답 지연을 통해 속도 조절
필요 이유
- 특정 프로듀서가 갑자기 대량의 메시지를 쏟아내면, 브로커의 네트워크/디스크 자원을 독점하게 되어 다른 프로듀서, 컨슈머의 성능이 저하될 수 있음
- 이를 막기 위해 브로커 관리자가 클라이언트별로 초당 처리 가능한 바이트 수(byte rate)를 제한
'공부 > Kafka' 카테고리의 다른 글
| [Kafka] Kafka CLI (0) | 2025.08.29 |
|---|---|
| [Kafka] Topic, Partition, Producer, Consumer (0) | 2025.08.29 |
| [Kafka] Zookeeper와 Broker (0) | 2025.06.29 |
| [Kafka] Kafka의 특징 (0) | 2021.06.13 |
| [Kafka: 카프카] Kafka란? (0) | 2021.06.07 |



