[Kafka] Producer와 메시지 전달 보장
개요
Kafka 프로듀서는 메시지를 직렬화하고, 어느 파티션으로 보낼지 결정하고, 배치로 묶고, 압축하고, 응답을 기다리고, 재시도하고, 중복을 막고, 트랜잭션을 묶는다. 브로커는 로그에 append만 잘하면 되는 단순한 서버이며, 신뢰성과 성능을 결정짓는 거의 모든 의사 결정이 프로듀서 쪽에서 일어난다.
Producer 내부 파이프라인
프로듀서가 send(record)를 호출하면 메시지는 다음 5단계를 거쳐 브로커로 향한다.
1단계: Serializer
ProducerRecord의 Key와 Value를 바이트 배열로 변환한다. key.serializer, value.serializer 설정으로 직렬화 클래스를 지정한다. StringSerializer, ByteArraySerializer, JsonSerializer, Avro/Protobuf 기반의 Confluent Serializer 등이 있다.
2단계: Partitioner
메시지가 토픽의 어느 파티션에 저장될지를 결정한다.
- 키가 있는 경우:
murmur2(key) % partition_count로 결정. 같은 키는 항상 같은 파티션으로 간다. - 키가 없는 경우 (Kafka 2.4+): Sticky Partitioner가 적용된다. 하나의 파티션에 배치가 찰 때까지 몰아 보내고, 가득 차면 다음 파티션으로 이동한다. 라운드로빈 방식이 배치를 작게 만들었던 문제를 해결한다.
- 커스텀 파티셔너:
partitioner.class로 지정할 수 있다. VIP 고객 전용 파티션, 지역 기반 파티셔닝 등에 활용한다.
3단계: Compression
compression.type 설정에 따라 메시지를 압축한다. none, gzip, snappy, lz4, zstd 중 선택한다. 압축은 배치 단위로 수행되므로 배치가 클수록 압축률이 좋다. 브로커는 압축 형태 그대로 저장하고, 컨슈머가 해제한다.
4단계: RecordAccumulator (배치 버퍼)
압축된 메시지는 즉시 전송되지 않고, 파티션별로 구분된 Deque<ProducerBatch>에 쌓인다. 이게 처리량과 지연 시간을 결정하는 핵심 단계다.
batch.size(기본 16KB): 배치 최대 크기. 가득 차면 즉시 전송 대상.linger.ms(기본 0): 배치가 가득 차지 않아도 이 시간이 지나면 전송.buffer.memory(기본 32MB): 전체 Accumulator의 총 메모리. 초과 시send()가max.block.ms(기본 60초) 동안 블로킹된다.
5단계: Sender Thread
사용자 스레드와는 별도로 백그라운드에서 동작하는 Sender 스레드가 Accumulator에서 준비된 배치를 가져와 브로커로 비동기 전송한다.
여기서 중요한 최적화가 피기배킹(piggybacking)이다. 아직 배치 사이즈 상한에 도달하지 못한 파티션 배치일지라도, 네트워크 통신 횟수를 줄이기 위해 동일한 브로커가 담당하는 파티션들이라면 한번에 묶어서 전송한다.
max.in.flight.requests.per.connection (기본 5)은 응답을 기다리지 않고 한 연결에 동시에 보낼 수 있는 최대 요청 수를 제한한다. 이 값이 1보다 크면 처리량은 올라가지만, 재시도 시 순서가 역전될 가능성이 생긴다.
핵심 설정 파라미터
acks: 응답을 어디까지 기다릴 것인가
acks는 프로듀서가 브로커로부터 어느 시점에 응답을 받을지를 결정한다.
sequenceDiagram participant P as Producer participant L as Leader participant F1 as Follower 1 (ISR) participant F2 as Follower 2 (ISR) rect rgb(232, 245, 233) Note over P,F2: acks=0 (Fire and Forget) P->>L: send(record) Note over P: 응답 안 기다림, 다음 메시지 송신 end rect rgb(255, 243, 224) Note over P,F2: acks=1 (Leader Only) P->>L: send(record) L->>L: 로그에 append L-->>P: ack L->>F1: 비동기 복제 L->>F2: 비동기 복제 end rect rgb(227, 242, 253) Note over P,F2: acks=all (All ISR) P->>L: send(record) L->>L: 로그에 append F1->>L: fetch L-->>F1: replicate F2->>L: fetch L-->>F2: replicate Note over L: 모든 ISR 복제 완료 후 High Watermark 이동 L-->>P: ack end
| 설정 | 응답 시점 | 손실 위험 | 처리량 | 권장 용도 |
|---|---|---|---|---|
acks=0 | 응답 안 기다림 | 가장 높음 | 가장 빠름 | 메트릭, 로그 샘플링 |
acks=1 | 리더만 받음 | 리더 장애 시 가능 | 중간 | 일반 이벤트 |
acks=all (-1) | 모든 ISR 복제 완료 | 거의 없음 | 가장 느림 | 금융, 주문 등 정합성 중요 |
acks=all은 브로커의 min.insync.replicas와 함께 써야 의미가 있다. ISR이 리더 하나뿐인 상황에서도 acks=all이 통과되면 사실상 acks=1과 같기 때문이다. 보통 replication-factor=3, min.insync.replicas=2로 두면 한 대가 죽어도 안전하게 동작한다.
batch.size, linger.ms: 처리량과 지연 트레이드오프
linger.ms를 0이 아닌 값으로 두면 약간의 지연을 감수하고 처리량을 크게 올릴 수 있다. 로그 수집처럼 처리량이 중요한 경우 linger.ms=20~50, batch.size=64KB~128KB 정도가 흔한 조합이다. 반대로 실시간 명령 전달처럼 지연이 중요한 경우 linger.ms=0을 그대로 둔다.
메시지 전달 보장 수준
분산 시스템에서 메시지 전달 보장은 세 가지로 나뉜다.
| 보장 수준 | 손실 | 중복 | 처리량 | 사용 사례 |
|---|---|---|---|---|
| At-Most-Once | 가능 | 없음 | 가장 빠름 | 실시간 로그, 메트릭 |
| At-Least-Once | 없음 | 가능 | 중간 | 일반 이벤트 (컨슈머가 멱등) |
| Exactly-Once | 없음 | 없음 | 가장 느림 | 금융 거래, 결제, 재고 |
At-Most-Once: 최대 한 번
응답을 기다리지 않고 보낸다. 빠르지만 네트워크 장애나 브로커 다운 시 손실 가능.
acks=0
retries=0실시간 로그 수집, 일부 손실이 허용되는 메트릭 데이터, IoT 센서 데이터 등에 적합하다.
At-Least-Once: 최소 한 번
리더의 ack을 기다리고, 못 받으면 재전송한다. 손실은 없지만 중복은 발생할 수 있다.
acks=all
retries=2147483647
max.in.flight.requests.per.connection=5중복이 발생하는 시나리오는 다음과 같다.
sequenceDiagram participant P as Producer participant L as Leader P->>L: 메시지 전송 (Seq=10) L->>L: 저장 성공 L--xP: ack 전송 실패 (네트워크 오류) Note over P: 타임아웃 → 재전송 P->>L: 동일 메시지 재전송 L->>L: 또 저장 (중복!) L-->>P: ack Note over L: 결과: 같은 메시지 2건
브로커는 메시지를 정상 저장했지만, 네트워크 오류로 응답이 프로듀서에 도달하지 못한 경우다. 프로듀서는 재전송하고, 결과적으로 중복 저장된다. 일반적인 이벤트 처리에서는 컨슈머 쪽에서 멱등성을 보장해 중복을 흡수하는 것이 표준 패턴이다.
Exactly-Once: 정확히 한 번
손실도 중복도 없는 보장은 세 가지 메커니즘의 조합으로 만들어진다.
acks=all: 모든 ISR에 복제 완료 후 응답으로 손실 방지.- 멱등성 프로듀서 (
enable.idempotence=true): PID + Sequence Number로 중복 자동 제거. - 트랜잭션 (
transactional.id설정): 여러 파티션·토픽에 걸친 원자적 발행.
이 세 가지가 모두 켜져야 Kafka 내부에서의 Exactly-Once Semantics(EOS)가 완성된다.
멱등성 프로듀서 (Idempotent Producer)
동작 원리
멱등성 프로듀서는 프로듀서 재시도로 인한 중복을 브로커 쪽에서 감지하고 제거한다. 각 메시지에 (PID, Sequence Number)라는 고유 식별자를 부여하는 것이 핵심이다.
sequenceDiagram participant P as Producer (PID=5, Epoch=0) participant B as Broker P->>B: send(Seq=0) B->>B: 저장 (last_seq=0) B-->>P: ack P->>B: send(Seq=1) B->>B: 저장 (last_seq=1) B-->>P: ack P->>B: send(Seq=2) B->>B: 저장 (last_seq=2) B--xP: ack 실패 (네트워크) Note over P: 타임아웃 → 재전송 P->>B: send(Seq=2) 재전송 Note over B: 중복 감지! (last_seq=2 == 2) B-->>P: ack (저장 안 함) Note over B: 결과: 메시지 1개만 존재
- PID (Producer ID): 프로듀서 인스턴스 고유 식별자. 브로커가 할당.
- Epoch: PID의 세대 번호. 재시작 시 증가.
- Sequence Number: 파티션별로 단조 증가하는 시퀀스. 0부터 시작.
브로커는 (PID, Partition) 별로 최근 5개 배치의 시퀀스 넘버를 메모리에 추적한다. 새로 도착한 메시지의 Seq가 이미 기록된 값보다 작거나 같으면 중복으로 판단하고 저장하지 않고 ack만 돌려준다.
설정
enable.idempotence=true
# 멱등성 활성화 시 자동으로 적용되는 값들:
acks=all
retries=2147483647
max.in.flight.requests.per.connection=5한계
멱등성 프로듀서가 보장하는 것은 단일 프로듀서 세션 내에서, 단일 파티션 내에서의 중복 제거다.
- 프로듀서가 재시작하면 새로운 PID가 할당되므로 이전 세션의 중복은 감지하지 못한다.
- 여러 파티션에 걸친 원자성은 보장하지 않는다 (예: 같은 트랜잭션에서 두 토픽에 발행하는 경우 한 토픽만 성공하고 다른 토픽이 실패할 수 있다).
트랜잭셔널 프로듀서
트랜잭션이란
트랜잭션은 여러 파티션·토픽에 대한 메시지 발행을 원자적(atomic)인 하나의 단위로 묶는다. All-or-Nothing 방식이며, 모두 커밋되거나 모두 abort된다. 또한 프로듀서 재시작 후에도 같은 transactional.id를 사용하면 멱등성이 세션을 넘어서 유지된다.
sequenceDiagram participant P as Producer participant TC as Transaction Coordinator participant B1 as Broker (Topic A) participant B2 as Broker (Topic B) participant C as Consumer Note over P: 1. 트랜잭션 초기화 P->>TC: initTransactions() TC->>P: PID + Epoch 할당<br/>(이전 미완료 트랜잭션 abort) Note over P: 2. 트랜잭션 시작 P->>TC: beginTransaction() Note over P: 3. 메시지 발행 (UNCOMMITTED) P->>B1: send(Topic A, msg1) P->>B2: send(Topic B, msg2) P->>B1: send(Topic A, msg3) Note over B1,B2: 저장되지만 UNCOMMITTED 상태 Note over P: 4. 커밋 P->>TC: commitTransaction() TC->>B1: COMMIT Marker 기록 TC->>B2: COMMIT Marker 기록 TC-->>P: 완료 Note over B1,B2: 메시지 COMMITTED 상태로 전환 Note over C: 5. read_committed 컨슈머 읽기 C->>B1: fetch (isolation.level=read_committed) B1-->>C: COMMITTED 메시지만 반환
주요 구성 요소
- Transaction ID (
transactional.id): 프로듀서를 식별하는 고유 ID. 재시작 후에도 같은 ID를 사용하면 코디네이터가 이전 미완료 트랜잭션을 자동으로 abort한다. - Transaction Coordinator: 트랜잭션 상태를 관리하는 브로커. 내부 토픽
__transaction_state에 상태를 저장한다. - Transaction Markers: COMMIT 또는 ABORT 마커. 각 파티션 로그에 기록되며, 컨슈머는 이 마커를 보고 메시지를 처리할지 결정한다.
Spring Kafka 트랜잭션 설정
@Configuration
public class KafkaTransactionConfig {
@Bean
public ProducerFactory<String, Object> producerFactory() {
Map<String, Object> config = new HashMap<>();
config.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
config.put(ProducerConfig.ACKS_CONFIG, "all");
config.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
config.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
config.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class);
DefaultKafkaProducerFactory<String, Object> factory =
new DefaultKafkaProducerFactory<>(config);
factory.setTransactionIdPrefix("order-tx-"); // 핵심: prefix 설정
return factory;
}
@Bean
public KafkaTemplate<String, Object> kafkaTemplate(ProducerFactory<String, Object> pf) {
return new KafkaTemplate<>(pf);
}
@Bean
public KafkaTransactionManager<String, Object> kafkaTransactionManager(
ProducerFactory<String, Object> pf) {
return new KafkaTransactionManager<>(pf);
}
}setTransactionIdPrefix()를 호출하면 Spring Kafka가 트랜잭셔널 프로듀서를 생성한다. KafkaTransactionManager를 빈으로 등록하면 @Transactional로 묶어 사용할 수 있다.
트랜잭션 사용 패턴
@Service
@RequiredArgsConstructor
public class OrderEventProducer {
private final KafkaTemplate<String, Object> kafkaTemplate;
// 패턴 1: executeInTransaction 람다
public void sendOrderAndAnalytics(OrderEvent event) {
kafkaTemplate.executeInTransaction(ops -> {
ops.send("orders", event.orderId(), event);
ops.send("order-analytics", event.orderId(), event);
ops.send("order-notifications", event.orderId(), event);
return true;
// 예외 발생 시 모든 메시지 abort
});
}
// 패턴 2: @Transactional 어노테이션
@Transactional("kafkaTransactionManager")
public void sendWithAnnotation(OrderEvent event) {
kafkaTemplate.send("orders", event.orderId(), event);
kafkaTemplate.send("order-analytics", event.orderId(), event);
// 메서드 정상 종료 → commit, 예외 → abort
}
}컨슈머 측 isolation.level
트랜잭셔널 프로듀서가 발행한 메시지를 정확하게 읽으려면 컨슈머에서 isolation.level=read_committed를 설정해야 한다.
# 컨슈머 설정
isolation.level=read_committed # COMMITTED 메시지만 읽음
# isolation.level=read_uncommitted (기본): 트랜잭션 무시read_committed로 설정하면 COMMITTED 상태가 될 때까지 해당 메시지는 컨슈머에게 보이지 않는다. 트랜잭션이 abort되면 해당 메시지는 영원히 건너뛴다(컨슈머는 보지 못함).
메시지 전달 보장 비교표
| 설정 | At-Most-Once | At-Least-Once | Exactly-Once |
|---|---|---|---|
acks | 0 | 1 | all |
enable.idempotence | false | false | true |
transactional.id | (미설정) | (미설정) | 설정 (선택) |
retries | 0 | > 0 | Integer.MAX_VALUE |
max.in.flight | 5 | 5 | 5 |
| 메시지 손실 | 가능 | 없음 | 없음 |
| 메시지 중복 | 없음 | 가능 | 없음 |
| 지연시간 | 낮음 | 중간 | 높음 |
| 처리량 | 높음 | 중간 | 낮음 |
파티셔닝 전략
키 기반 파티셔닝의 주의점
키가 있으면 murmur2(key) % partition_count로 파티션이 결정되어 같은 키의 메시지는 항상 같은 파티션으로 간다.
그런데 한 가지 주의점이 있다. 파티션 수를 변경하면 키-파티션 매핑이 깨진다. 기존에는 userId=42가 파티션 3으로 갔다면, 파티션을 12개에서 24개로 늘리면 userId=42가 갑자기 파티션 18로 가게 될 수 있다. 결과적으로 같은 키의 메시지가 서로 다른 파티션에 흩어지면서 순서 보장이 깨진다.
따라서 파티션 수는 초기에 충분히 잡아두고, 부하 측정 후 점진적으로 늘린다.
커스텀 파티셔너
public class VipPartitioner implements Partitioner {
@Override
public int partition(String topic, Object key, byte[] keyBytes,
Object value, byte[] valueBytes, Cluster cluster) {
int partitionCount = cluster.partitionCountForTopic(topic);
if (value instanceof OrderEvent order && order.isVip()) {
return 0; // VIP는 항상 파티션 0
}
return Math.abs(key.hashCode() % (partitionCount - 1)) + 1;
}
@Override public void configure(Map<String, ?> configs) {}
@Override public void close() {}
}partitioner.class=com.example.VipPartitioner코드 예시
Spring Kafka 기본 프로듀서
@Configuration
public class KafkaProducerConfig {
@Bean
public ProducerFactory<String, Object> producerFactory() {
Map<String, Object> config = Map.of(
ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092",
ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class,
ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class,
ProducerConfig.ACKS_CONFIG, "all",
ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true,
ProducerConfig.COMPRESSION_TYPE_CONFIG, "lz4",
ProducerConfig.LINGER_MS_CONFIG, 20,
ProducerConfig.BATCH_SIZE_CONFIG, 32768
);
return new DefaultKafkaProducerFactory<>(config);
}
@Bean
public KafkaTemplate<String, Object> kafkaTemplate(ProducerFactory<String, Object> pf) {
return new KafkaTemplate<>(pf);
}
}동기/비동기 전송
@Service
@RequiredArgsConstructor
@Slf4j
public class OrderEventProducer {
private final KafkaTemplate<String, Object> kafkaTemplate;
// 동기 전송 (응답을 기다림, 성능 저하 주의)
public void sendSync(String topic, OrderEvent event) {
try {
SendResult<String, Object> result =
kafkaTemplate.send(topic, event.orderId(), event)
.get(10, TimeUnit.SECONDS);
RecordMetadata metadata = result.getRecordMetadata();
log.info("Sent: topic={}, partition={}, offset={}",
metadata.topic(), metadata.partition(), metadata.offset());
} catch (Exception e) {
log.error("Send failed", e);
throw new KafkaSendException(e);
}
}
// 비동기 전송 (권장, Sender 스레드 활용 극대화)
public void sendAsync(String topic, OrderEvent event) {
CompletableFuture<SendResult<String, Object>> future =
kafkaTemplate.send(topic, event.orderId(), event);
future.whenComplete((result, ex) -> {
if (ex != null) {
log.error("Send failed: {}", ex.getMessage(), ex);
// 실패 처리: 보상 트랜잭션, 알림, DLQ 등
} else {
log.debug("Sent partition={} offset={}",
result.getRecordMetadata().partition(),
result.getRecordMetadata().offset());
}
});
}
// 헤더와 함께 전송 (트레이싱, 이벤트 타입 등)
public void sendWithHeaders(OrderEvent event) {
ProducerRecord<String, Object> record = new ProducerRecord<>(
"orders", null, event.orderId(), event);
record.headers()
.add("event-type", "ORDER_CREATED".getBytes())
.add("source", "order-service".getBytes())
.add("trace-id", MDC.get("traceId").getBytes());
kafkaTemplate.send(record);
}
}Kotlin (Spring Boot 3.x)
@Service
class PaymentEventProducer(
private val kafkaTemplate: KafkaTemplate<String, Any>,
) {
private val log = LoggerFactory.getLogger(javaClass)
fun publish(payment: PaymentEvent) {
kafkaTemplate.send("payment-events", payment.paymentId, payment)
.whenComplete { result, ex ->
if (ex != null) {
log.error("publish failed: ${payment.paymentId}", ex)
} else {
log.debug(
"published: partition={}, offset={}",
result.recordMetadata.partition(),
result.recordMetadata.offset(),
)
}
}
}
}application.yml 용도별 권장 설정
# 높은 신뢰성 (금융, 주문)
spring:
kafka:
producer:
acks: all
retries: 2147483647
properties:
enable.idempotence: true
max.in.flight.requests.per.connection: 5
delivery.timeout.ms: 120000
# 브로커: min.insync.replicas=2
---
# 높은 처리량 (로그, 메트릭)
spring:
kafka:
producer:
acks: 1
properties:
batch.size: 65536
linger.ms: 20
compression.type: lz4
buffer.memory: 67108864
---
# 낮은 지연 (실시간 명령)
spring:
kafka:
producer:
acks: 1
properties:
linger.ms: 0
batch.size: 16384
compression.type: nonekafka-cli 명령
# 트랜잭션 상태 토픽 확인
kafka-topics.sh --bootstrap-server localhost:9092 \
--describe --topic __transaction_state
# 콘솔 프로듀서로 메시지 발행 (테스트용)
kafka-console-producer.sh --bootstrap-server localhost:9092 \
--topic orders \
--property "parse.key=true" --property "key.separator=:" \
--property "acks=all"
# 컨슈머에서 read_committed로 트랜잭션 메시지만 보기
kafka-console-consumer.sh --bootstrap-server localhost:9092 \
--topic orders --from-beginning \
--isolation-level read_committed