티스토리 뷰
아래 글은 Apache Kafka Producer의 핵심을 “개념 → 주요 설정 → Java/Python 예제 → 성능·신뢰성 모범 사례” 순으로 정리한 전문가용 가이드입니다. 프로듀서가 하는 일, acks·직렬화·배치·압축·Idempotence 등 필수 설정의 의미와 실무 권장값을 설명하고, 즉시 실행 가능한 최소-예제 코드를 제공합니다. 이 글을 따라 하면 “Why & How” 를 모두 이해하고 안전-고성능으로 메시지를 발행할 수 있습니다.
1. 프로듀서란 무엇인가?
Kafka Producer는 사용자가 작성한 이벤트를 파티션 단위 로그에 기록하여 브로커에 전달하는 클라이언트 컴포넌트입니다. 프로듀서는 레코드를 직렬화 → 파티셔닝 → 배치·압축 → 전송 → 재시도/확인(acks) 과정을 거치며, 내부 버퍼와 전송 쓰레드가 비동기로 동작해 TPS를 극대화합니다.
2. 주요 설정 한눈에 보기
카테고리 설정 키 의미-특징 실무 권장
| 내구성 | acks | 0: 무확인, 1: 리더 확인, all/-1: ISR 전원 확인 | all (거의 필수) |
| retries / delivery.timeout.ms | 전송 실패 재시도 횟수·총 대기시간 | retries=Integer.MAX_VALUE, delivery.timeout.ms 충분히 크게 | |
| 정확-한-번 | enable.idempotence | 중복 방지·순서 유지(Seq no) | true → 자동으로 acks=all·retries 조정 |
| 배치 | batch.size | 파티션별 메모리 버퍼(byte) | 네트워크 MTU 맞춰 32 KB~128 KB |
| linger.ms | 최대 대기 시간(ms) | 1-10 ms로 소폭 지연해 TPS↑ | |
| 압축 | compression.type | snappy / lz4 / zstd | 네트워크 IO 지연이 큰 환경은 lz4 추천 |
| 직렬화 | key.serializer / value.serializer | 바이트 변환 클래스 | StringSerializer, AvroSerializer 등 |
| 파티셔닝 | partitioner.class | 사용자 정의 파티션 규칙 | 해시 불균형 시 커스텀 구현 |
팁: enable.idempotence=true 로두면 max.in.flight.requests.per.connection 이 5 이하로 자동 제한돼 순서 역전이 발생하지 않습니다.
3. Java 프로듀서 기본 예제
import org.apache.kafka.clients.producer.*;
import java.util.Properties;
public class QuickStartProducer {
public static void main(String[] args) throws Exception {
Properties p = new Properties();
p.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
p.put(ProducerConfig.ACKS_CONFIG, "all");
p.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, "true");
p.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,
"org.apache.kafka.common.serialization.StringSerializer");
p.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,
"org.apache.kafka.common.serialization.StringSerializer");
p.put(ProducerConfig.LINGER_MS_CONFIG, 5);
p.put(ProducerConfig.BATCH_SIZE_CONFIG, 64_000); // 64 KB
try (KafkaProducer<String, String> producer = new KafkaProducer<>(p)) {
ProducerRecord<String,String> rec =
new ProducerRecord<>("quickstart", "key1", "안녕하세요 Kafka!");
RecordMetadata meta = producer.send(rec).get(); // 동기 – 데모용
System.out.printf("partition=%d offset=%d%n", meta.partition(), meta.offset());
}
}
}
위 코드는 acks=all+idempotence 로 정확-한-번 전송을 만족하며, linger.ms·batch.size 로 미세 배치를 활성화했습니다.
4. Python (kafka-python) 기본 예제
from kafka import KafkaProducer
import json, time
producer = KafkaProducer(
bootstrap_servers='localhost:9092',
acks='all',
retries=2147483647,
linger_ms=5,
batch_size=64_000,
key_serializer=str.encode,
value_serializer=lambda v: json.dumps(v).encode(),
enable_idempotence=True # kafka-python 2.0+
)
payload = {"orderId": 123, "amount": 19.99}
producer.send('quickstart', key='order-123', value=payload)
producer.flush()
producer.close(timeout=5)
kafka-python 의 acks 파라미터는 0/1/all 값을 동일하게 지원하며, 직렬화 함수를 람다로 지정하여 JSON 인코딩을 수행합니다.
5. 성능·신뢰성 모범 사례
5-1. Throughput 튜닝
- 배치 & 지연: linger.ms 를 1-10 ms 로 설정해 자연스러운 배치 형성 후 전송량을 2-3 배까지 끌어올릴 수 있습니다.
- 압축: 레코드가 수 KB 이상이거나 네트워크 지연이 긴 경우 compression.type=lz4 가 CPU 대비 전송량 효율이 우수합니다.
5-2. 정확-한-번 보장
- enable.idempotence=true + acks=all + retries 무제한이면 단일 파티션에서 “at-least-once → exactly-once” 로 승격됩니다.
- 멀티 파티션 트랜잭션은 transactional.id 지정 후 initTransactions / beginTransaction / commitTransaction API를 호출해 구현합니다.
5-3. 오류 처리
- onCompletion 콜백 또는 Future#get 예외를 확인해 TimeoutException, RecordTooLargeException 등을 로그로 남기고 비즈니스 재시도 정책을 세웁니다.
- 재시도 시 동일 키를 사용하면 파티션·순서가 유지되어 컨슈머에서 중복 필터링이 수월합니다.
5-4. 스키마 관리
- Avro/Protobuf/JSON Schema 를 쓰는 경우 Schema Registry 와 사전-호환성 정책(BACKWARD/FORWARD) 으로 데이터 진화 비용을 최소화하세요.
6. 마무리
Kafka Producer는 단순 “send” API 뒤에 내구성·순서·배치·압축 을 절묘하게 조율하는 고성능 엔진입니다. acks=all + Idempotence 로 안전을 우선 확보한 뒤, linger.ms·batch.size·compression.type 을 조정해 필요 TPS·지연 목표를 달성하세요. Java/Python 예제 코드를 기반으로 자신만의 프로듀서 모듈을 패키징해 두면, 데이터 파이프라인 전체 품질이 한 단계 올라갑니다. Happy Streaming!
'개발 인프라 > 카프카' 카테고리의 다른 글
| 주키퍼(ZooKeeper)는 이제 안녕? KRaft 모드 알아보기 (0) | 2025.05.22 |
|---|---|
| 카프카 토픽과 파티션, 왜 중요할까? (1) | 2025.05.21 |
| 카프카 컨슈머(Consumer)와 컨슈머 그룹 파헤치기 (1) | 2025.05.16 |
| 로컬 카프카 설치부터 메시지 전송/수신 (1) | 2025.05.16 |
| 카프카(Kafka)란 무엇인가? (0) | 2025.05.14 |
- Total
- Today
- Yesterday
- 스브링부트
- Subagent
- 카프카 개념
- Stack Area
- MCP
- First-class citizen
- generated_body()
- unreal engjin
- vite
- 자바
- model context protocol
- JVM
- RESTfull
- redis
- Claude Agent SDK
- AI 에이전트
- Java
- ai통합
- JAVA 프로그래밍
- method Area
- Heap Area
- cqrs
- 언리얼엔진5
- 일급 객체
- 코틀린
- 타입 안전성
- 언리얼엔진
- 디자인패턴
- springai
- 코프링
| 일 | 월 | 화 | 수 | 목 | 금 | 토 |
|---|---|---|---|---|---|---|
| 1 | 2 | 3 | 4 | |||
| 5 | 6 | 7 | 8 | 9 | 10 | 11 |
| 12 | 13 | 14 | 15 | 16 | 17 | 18 |
| 19 | 20 | 21 | 22 | 23 | 24 | 25 |
| 26 | 27 | 28 | 29 | 30 | 31 |
