전체 그래프
Kafka

Kafka Producer

data-engineeringtool/kafkaconcept/producer

상위: Kafka

요약

Producer는 Kafka Topic으로 데이터를 전송하는 클라이언트입니다. 다양한 데이터 소스에서 메시지를 생성하여 Kafka 브로커로 전송하며, 직렬화, 파티셔닝, 압축 등의 기능을 제공합니다.

메시지 전송 과정

image 60 1.png

  1. 직렬화 (Serialization): 데이터를 Byte 형태로 변환
  2. 파티션 결정 (Partitioner): 어느 파티션으로 보낼지 결정
  3. 압축 (Compression): 데이터 압축 (선택)
  4. RA 버퍼 (Record Accumulator): 배치 단위로 모음
  5. 전송: Sender Thread가 브로커로 전송

Producer Record

메시지의 기본 단위

image 61 1.png

필드필수 여부설명
Topic필수전송될 Kafka 토픽
Value필수전송할 데이터 값
Key선택파티션을 지정할 키 값
Partition선택전송될 파티션 번호
Headers선택포함할 기타 정보
Timestamp선택생성 시간

직렬화 (Serialization)

image 62 1.png

장점

  1. 네트워크 전송 유리
  2. 저장 빠름
  3. 압축 용이
  4. 복구 용이

주요 Serializer

  • StringSerializer: 단순 문자열
  • ByteArraySerializer: 바이트 배열
  • JsonSerializer: JSON 포맷
  • AvroSerializer: 스키마 기반 (Kafka 권장)
  • ProtobufSerializer: 최고 압축률

파티셔닝 전략

1. Round Robin

파티션 지정이 없을 경우 순차적으로 분배

image 63 1.png

2. Key Base

Key가 있으면 해시를 통해 같은 Key는 같은 파티션에 저장

image 64 1.png

image 65 1.png

주의: 파티션 수 변경 시 Key 매핑 변경됨

3. 파티션 직접 지정

image 66 1.png

image 67 1.png

4. Uniform Sticky (Default)

하나의 목표 파티션을 빠르게 채우고 목표 파티션을 바꿈

image 68 1.png

image 69 1.png

압축 (Compression)

image 70 1.png

압축 방식압축률속도CPU 사용용도
Gzip높음느림높음확실한 압축
LZ4중간중간중간균형잡힌 선택
Snappy낮음빠름낮음빠른 압축

Batching

RA (Record Accumulator) 버퍼

image 71 1.png

image 72 1.png

image 73 1.png

주요 설정

  • batch.size: 하나의 배치 크기 (기본값: 16KB)
  • linger.ms: 배치 대기 시간 (기본값: 0ms)
  • buffer.memory: 전체 버퍼 메모리 (기본값: 32MB)

PiggyBack

조건을 만족하지 않은 배치도, 만족한 배치와 같은 브로커일 경우 함께 전송

image 74 1.png

Acknowledge 옵션

image 77 1.png

image 78 1.png

image 79 1.png

  • acks=0: 브로커 수신 여부 확인 안 함
  • acks=1: 리더가 수신하면 성공
  • acks=all(-1): 리더 + 팔로워 모두 복제 완료 시 성공 (기본값)

min.insync.replicas

image 80 1.png

ISR에 포함된 follower 중 최소 몇 개가 복제를 완료해야 ack를 주는지 설정

Idempotence (멱등성)

image 81 1.png

  • 중복 전송 방지 기능
  • 같은 메시지가 중복 전송되어도 중복 저장되지 않음
  • 현재 Kafka에서 기본적으로 활성화됨
  • Exactly-Once 보장

Python 예제

from kafka import KafkaProducer

producer = KafkaProducer(
    bootstrap_servers='localhost:9092',
    value_serializer=lambda v: v.encode('utf-8'),
    acks='all',
    compression_type='lz4',
    batch_size=32768,
    linger_ms=5
)

# 메시지 전송
for i in range(5):
    message = f"hello kafka {i}"
    producer.send('test-topic', message)
    print(f"Sent: {message}")

producer.flush()
producer.close()