실무 예제 · 3
안전한 Producer 설정 (무손실)
"유실이 없다"는 프로듀서 설정 하나로 만들어지지 않습니다.
프로듀서의 acks·멱등성, 토픽의 min.insync.replicas,
브로커의 복제 계수 세 축이 동시에 맞아야 성립합니다.
이 예제는 세 축을 모두 갖춘 상태에서
실패를 조용히 삼키지 않는 프로듀서를 처음부터 끝까지 작성합니다.
학습 목표
acks=all이 실제로 무엇을 기다리는지,min.insync.replicas와 어떻게 결합해 내구성을 만드는지 설명할 수 있습니다.- 멱등 프로듀서가 중복을 제거하는 범위와 한계를 구분할 수 있습니다.
retries가 아니라delivery.timeout.ms가 재시도의 실질 상한임을 코드로 확인할 수 있습니다.- 콜백에서 retriable / non-retriable 예외를 분류해 유실 없이 실패를 보고하는 코드를 작성할 수 있습니다.
시나리오
결제 게이트웨이가 카드사 승인 결과를 받으면 orders 토픽에
OrderApproved 이벤트를 발행합니다. 이 이벤트를 정산 배치와 알림 서비스가 소비합니다.
하루 40만 건, 피크 시 초당 120건 수준입니다.
문제는 이 이벤트가 없으면 정산이 누락된다는 점입니다. 고객은 이미 결제했는데 정산 데이터가 없으면 회계상 미수금이 되고, 발견되는 시점은 월말 마감입니다. 따라서 요구사항은 "쓰기가 실패했다면 반드시 실패했다는 사실을 알아야 한다"입니다. 느려지는 것은 허용하지만 조용히 사라지는 것은 허용하지 않습니다.
아키텍처
브로커 3대, 복제 계수 3, min.insync.replicas=2입니다.
acks=all은 "모든 레플리카"가 아니라
"현재 ISR에 있는 레플리카 전부"를 기다립니다.
그래서 ISR이 1로 쪼그라들면 acks=all만으로는 사본이 하나뿐인 쓰기를 성공으로 인정하게 됩니다.
min.insync.replicas=2가 그 지점에서
성공 대신 실패를 반환하게 만드는 안전장치입니다.
acks 0 / 1 / all 시퀀스 비교 — 각 설정에서 유실이 발생하는 지점
min.insync.replicas에 미달하는 과정
| 조합 | 브로커 1대 정지 | 브로커 2대 정지 | 유실 가능성 |
|---|---|---|---|
acks=0 |
쓰기 성공(응답을 안 봄) | 쓰기 성공 | 있음 — 브로커에 도달했는지조차 확인하지 않습니다 |
acks=1 |
쓰기 성공 | 쓰기 성공 | 있음 — 리더가 응답 후 복제 전에 죽으면 유실 |
acks=all + min.insync.replicas=1 |
쓰기 성공 | 쓰기 성공(사본 1개) | 있음 — 남은 1대가 죽으면 유실 |
acks=all + min.insync.replicas=2 |
쓰기 성공(사본 2개) | 쓰기 실패 | 없음 — 실패를 반환하므로 애플리케이션이 대응할 수 있습니다 |
acks=all + min.insync.replicas=3 |
쓰기 실패 | 쓰기 실패 | 없음 — 다만 가용성이 너무 낮습니다(1대 정비만으로 중단) |
사전 요구사항
- 예제 1의 3노드 KRaft 클러스터가 떠 있어야 합니다. 호스트에서
localhost:29092,localhost:39092,localhost:49092로 접속됩니다. - Apache Kafka 4.3.1 클라이언트 (
org.apache.kafka:kafka-clients) - Java 17 이상 (Kafka 4.3은 17 · 21 · 25를 완전 지원합니다)
- Maven 3.9 이상 또는 동등한 Gradle 설정
orders 토픽은 예제 1의 create-topics.sh가
파티션 6 · RF 3 · min.insync.replicas=2로 만들어 둡니다.
아래 코드에는 AdminClient로 같은 토픽을 보장하는 프로비저너도 포함했으니
토픽이 없어도 그대로 실행됩니다.
전체 코드
safe-producer/
├── pom.xml
└── src/main/java/com/example/kafka/
├── TopicProvisioner.java # 토픽을 요구 스펙대로 보장
├── SafeProducerConfig.java # 무손실 설정을 한곳에 모음
├── OrderApprovedProducer.java# 프로듀서 본체 (콜백·예외 분류·종료 처리)
└── Main.java # 진입점
pom.xml
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>com.example</groupId>
<artifactId>safe-producer</artifactId>
<version>1.0.0</version>
<packaging>jar</packaging>
<properties>
<maven.compiler.release>17</maven.compiler.release>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<!-- 브로커와 클라이언트 버전을 맞춥니다. 클라이언트가 브로커보다 낮아도 동작하지만
4.3 에 추가된 API 를 쓰려면 클라이언트도 4.3 이어야 합니다. -->
<kafka.version>4.3.1</kafka.version>
<slf4j.version>2.0.17</slf4j.version>
</properties>
<dependencies>
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>${kafka.version}</version>
</dependency>
<!-- kafka-clients 는 slf4j-api 만 의존합니다.
바인딩이 없으면 경고만 나오고 로그가 하나도 안 보입니다. -->
<dependency>
<groupId>org.slf4j</groupId>
<artifactId>slf4j-simple</artifactId>
<version>${slf4j.version}</version>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.codehaus.mojo</groupId>
<artifactId>exec-maven-plugin</artifactId>
<version>3.5.0</version>
<configuration>
<mainClass>com.example.kafka.Main</mainClass>
</configuration>
</plugin>
</plugins>
</build>
</project>
무손실 설정 — 값마다 이유가 있습니다
package com.example.kafka;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.common.serialization.StringSerializer;
import java.util.Properties;
/**
* 유실을 허용하지 않는 프로듀서 설정.
*
* Kafka 4.3 기준으로 acks=all, enable.idempotence=true 는 이미 기본값이지만
* "기본값이라 안 썼다" 와 "의도해서 이 값이다" 는 운영 중에 구분되지 않습니다.
* 내구성에 직결되는 값은 기본값과 같아도 명시합니다.
*/
public final class SafeProducerConfig {
private SafeProducerConfig() {
}
public static Properties create(String bootstrapServers, String clientId) {
Properties p = new Properties();
// --- 연결 -----------------------------------------------------------
// 초기 연결 이중화. 부하 분산 목적이 아니라 부트스트랩 단일 장애점 제거입니다.
p.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
// 브로커 로그와 JMX 메트릭에서 이 애플리케이션을 식별하는 이름.
// 장애 조사 때 "누가 보낸 요청인지" 를 알려 주므로 반드시 지정합니다.
p.put(ProducerConfig.CLIENT_ID_CONFIG, clientId);
p.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
p.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
// --- 내구성 (이 예제의 핵심) -----------------------------------------
// acks=all : 현재 ISR 에 있는 레플리카 전부가 기록을 마쳐야 성공.
// "모든 레플리카" 가 아니라 "ISR 전부" 라는 점이 중요합니다.
// 4.3 기본값도 all 이지만 의도를 남기기 위해 명시합니다.
// 토픽의 min.insync.replicas=2 와 짝을 이뤄야 실효가 있습니다.
p.put(ProducerConfig.ACKS_CONFIG, "all");
// 멱등성 : 브로커가 (PID, 파티션, 시퀀스번호) 로 중복을 판별해
// "재시도로 인한 중복" 을 제거합니다. 4.3 기본값 true.
// 켜면 acks=all, retries=Integer.MAX_VALUE,
// max.in.flight<=5 제약이 자동으로 따라옵니다.
p.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
// 멱등성이 켜진 상태에서 in-flight 요청 5개까지는 순서가 보장됩니다.
// 6 이상으로 올리면 ConfigException 이 발생해 기동 자체가 실패합니다.
// 기본값이 5이므로 그대로 두면 되지만, 성능 튜닝 중 실수로
// 올리는 일이 잦아 명시해 잠가 둡니다.
p.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, 5);
// retries 를 직접 만지지 않습니다.
// 멱등성이 켜지면 기본값이 Integer.MAX_VALUE 이고,
// 실질적인 재시도 상한은 아래 delivery.timeout.ms 입니다.
// retries=0 으로 두면 첫 실패에서 곧바로 포기해 유실 위험이 생깁니다.
// --- 시간 예산 ------------------------------------------------------
// send() 부터 성공/실패가 확정될 때까지의 총 상한. 기본값 120000(2분).
// 이 값 안에서 request.timeout.ms 단위의 시도가 반복됩니다.
// 관계: delivery.timeout.ms >= linger.ms + request.timeout.ms
p.put(ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG, 120_000);
// 한 번의 요청이 브로커 응답을 기다리는 시간. 기본값 30000.
p.put(ProducerConfig.REQUEST_TIMEOUT_MS_CONFIG, 30_000);
// 버퍼가 꽉 찼거나 메타데이터를 못 받았을 때 send() 가 블록되는 최대 시간.
// 기본값 60000. 이 시간이 지나면 TimeoutException 으로 실패합니다.
// 0 에 가깝게 줄이면 순간 스파이크에 전송이 실패하고,
// 너무 키우면 호출 스레드가 오래 묶입니다.
p.put(ProducerConfig.MAX_BLOCK_MS_CONFIG, 60_000);
// --- 처리량 / 지연 --------------------------------------------------
// 배치 최대 크기(바이트). 기본값 16384.
// 결제 이벤트는 건당 1KB 미만이므로 64KB 로 올려 배치 효율을 높입니다.
p.put(ProducerConfig.BATCH_SIZE_CONFIG, 64 * 1024);
// 배치가 다 차지 않아도 이만큼 기다립니다. 4.0 에서 기본값이 0 -> 5 로 바뀌었습니다.
// 결제 이벤트는 20ms 지연을 감수하고 배치를 키우는 편이 유리합니다.
p.put(ProducerConfig.LINGER_MS_CONFIG, 20);
// 전송 대기 버퍼 총량. 기본값 33554432(32MB).
// 이 버퍼가 고갈되면 send() 가 max.block.ms 만큼 블록됩니다.
p.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 64 * 1024 * 1024L);
// 압축. 기본값은 none 입니다.
// JSON 텍스트는 압축 효율이 좋고, lz4 는 CPU 대비 압축률 균형이 좋습니다.
// 주의: 브로커의 message.max.bytes(기본 1048588) 는 "압축 후" 배치 크기 기준입니다.
p.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "lz4");
return p;
}
}
토픽 보장 — 코드가 요구 스펙을 선언합니다
토픽이 min.insync.replicas=2로 만들어졌는지를
애플리케이션이 확인하게 만들면, 잘못된 토픽에 쓰기 시작하는 사고를 막을 수 있습니다.
자동 생성(auto.create.topics.enable)에 의존하면
파티션 1 · RF 1 토픽이 조용히 만들어집니다.
package com.example.kafka;
import org.apache.kafka.clients.admin.Admin;
import org.apache.kafka.clients.admin.AdminClientConfig;
import org.apache.kafka.clients.admin.NewTopic;
import org.apache.kafka.clients.admin.TopicDescription;
import org.apache.kafka.common.config.TopicConfig;
import org.apache.kafka.common.errors.TopicExistsException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.List;
import java.util.Map;
import java.util.Properties;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
/**
* 토픽이 요구 스펙(파티션 수 / RF / min.insync.replicas)으로 존재하는지 보장합니다.
* 이미 있으면 만들지 않고, 스펙이 다르면 경고를 남깁니다.
* (파티션 수와 RF 는 생성 후 변경 방식이 다르므로 자동으로 고치지 않습니다.)
*/
public final class TopicProvisioner {
private static final Logger log = LoggerFactory.getLogger(TopicProvisioner.class);
private static final int ADMIN_TIMEOUT_SEC = 30;
private TopicProvisioner() {
}
public static void ensureTopic(String bootstrapServers,
String topic,
int partitions,
short replicationFactor,
String minInsyncReplicas) {
Properties props = new Properties();
props.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
// 관리 요청도 무한정 기다리지 않게 상한을 둡니다.
props.put(AdminClientConfig.REQUEST_TIMEOUT_MS_CONFIG, 15_000);
props.put(AdminClientConfig.DEFAULT_API_TIMEOUT_MS_CONFIG, 30_000);
// try-with-resources 로 닫습니다. Admin 은 내부에 스레드를 갖습니다.
try (Admin admin = Admin.create(props)) {
NewTopic newTopic = new NewTopic(topic, partitions, replicationFactor)
.configs(Map.of(
// acks=all 과 짝을 이루는 내구성 하한.
// RF=3 에 2 가 표준 조합입니다.
TopicConfig.MIN_IN_SYNC_REPLICAS_CONFIG, minInsyncReplicas));
admin.createTopics(List.of(newTopic))
.all()
.get(ADMIN_TIMEOUT_SEC, TimeUnit.SECONDS);
log.info("토픽 생성 완료: {} (partitions={}, rf={}, min.insync.replicas={})",
topic, partitions, replicationFactor, minInsyncReplicas);
} catch (ExecutionException e) {
// 이미 존재하는 것은 정상 경로입니다. 그 외는 그대로 올립니다.
if (e.getCause() instanceof TopicExistsException) {
log.info("토픽이 이미 존재합니다: {} — 스펙을 검증합니다", topic);
verify(props, topic, partitions, replicationFactor);
} else {
throw new IllegalStateException("토픽 생성 실패: " + topic, e.getCause());
}
} catch (TimeoutException e) {
throw new IllegalStateException(
"토픽 생성이 " + ADMIN_TIMEOUT_SEC + "초 안에 끝나지 않았습니다: " + topic, e);
} catch (InterruptedException e) {
// 인터럽트 상태를 반드시 복원합니다. 삼키면 상위 종료 로직이 깨집니다.
Thread.currentThread().interrupt();
throw new IllegalStateException("토픽 생성 중 인터럽트: " + topic, e);
}
}
private static void verify(Properties props, String topic,
int expectedPartitions, short expectedRf) {
try (Admin admin = Admin.create(props)) {
TopicDescription desc = admin.describeTopics(List.of(topic))
.allTopicNames()
.get(ADMIN_TIMEOUT_SEC, TimeUnit.SECONDS)
.get(topic);
int actualPartitions = desc.partitions().size();
int actualRf = desc.partitions().get(0).replicas().size();
if (actualPartitions != expectedPartitions) {
log.warn("파티션 수 불일치: {} — 기대 {}, 실제 {}. "
+ "파티션은 늘릴 수만 있고 늘리면 키 분배가 바뀝니다.",
topic, expectedPartitions, actualPartitions);
}
if (actualRf != expectedRf) {
log.warn("복제 계수 불일치: {} — 기대 {}, 실제 {}. "
+ "RF 는 kafka-reassign-partitions.sh 로만 바꿀 수 있습니다.",
topic, expectedRf, actualRf);
}
} catch (ExecutionException | TimeoutException e) {
log.warn("토픽 스펙 검증 실패: {}", topic, e);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
log.warn("토픽 스펙 검증 중 인터럽트: {}", topic);
}
}
}
프로듀서 본체 — 실패를 절대 삼키지 않습니다
package com.example.kafka;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.clients.producer.RecordMetadata;
import org.apache.kafka.common.KafkaException;
import org.apache.kafka.common.errors.AuthorizationException;
import org.apache.kafka.common.errors.RecordTooLargeException;
import org.apache.kafka.common.errors.RetriableException;
import org.apache.kafka.common.errors.SerializationException;
import org.apache.kafka.common.errors.UnknownTopicOrPartitionException;
import org.apache.kafka.common.header.internals.RecordHeader;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.nio.charset.StandardCharsets;
import java.time.Duration;
import java.util.Properties;
import java.util.concurrent.atomic.AtomicLong;
/**
* 결제 승인 이벤트를 유실 없이 발행합니다.
*
* 유실 방지의 실제 조건은 네 가지가 동시에 만족될 때입니다.
* 1) acks=all + 토픽 min.insync.replicas=2 (브로커 측 내구성)
* 2) enable.idempotence=true (재시도 중복 제거)
* 3) 콜백에서 실패를 집계하고 보고 (애플리케이션 측 인지)
* 4) 종료 시 flush() 로 버퍼를 비움 (프로세스 종료 시 유실 방지)
*/
public class OrderApprovedProducer implements AutoCloseable {
private static final Logger log = LoggerFactory.getLogger(OrderApprovedProducer.class);
private final Producer<String, String> producer;
private final String topic;
// 콜백은 Sender 스레드에서 실행됩니다. 여러 스레드가 만지므로 원자적 카운터를 씁니다.
private final AtomicLong sent = new AtomicLong();
private final AtomicLong failed = new AtomicLong();
private final AtomicLong fatal = new AtomicLong();
public OrderApprovedProducer(String bootstrapServers, String topic, String clientId) {
this.topic = topic;
this.producer = new KafkaProducer<>(
SafeProducerConfig.create(bootstrapServers, clientId));
}
/**
* 이벤트 1건을 비동기로 보냅니다.
*
* @param orderId 파티션 키. 같은 주문의 이벤트가 같은 파티션으로 가야
* 순서가 보장됩니다. 키가 null 이면 순서 보장의 근거가 없습니다.
* @param payload 이벤트 본문(JSON 문자열)
*/
public void send(String orderId, String payload) {
ProducerRecord<String, String> record = new ProducerRecord<>(topic, orderId, payload);
// 헤더는 라우팅·추적용 메타데이터를 본문 파싱 없이 읽게 해 줍니다.
// 컨슈머가 역직렬화에 실패해도 헤더는 읽을 수 있습니다.
record.headers()
.add(new RecordHeader("event-type", "OrderApproved".getBytes(StandardCharsets.UTF_8)))
.add(new RecordHeader("schema-version", "1".getBytes(StandardCharsets.UTF_8)));
try {
// send() 는 즉시 반환합니다. 결과는 아래 콜백으로 옵니다.
// 콜백은 반드시 논블로킹이어야 합니다 — Sender 스레드를 막으면
// 그 프로듀서의 모든 전송이 함께 멈춥니다.
producer.send(record, (RecordMetadata metadata, Exception exception) -> {
if (exception == null) {
sent.incrementAndGet();
if (log.isDebugEnabled()) {
log.debug("전송 성공 orderId={} {}-{}@{}",
orderId, metadata.topic(), metadata.partition(), metadata.offset());
}
return;
}
handleSendFailure(orderId, payload, exception);
});
} catch (SerializationException e) {
// 직렬화 실패는 send() 호출 스레드에서 즉시 던져집니다.
// 재시도해도 결과가 같으므로 데이터 자체를 격리해야 합니다.
fatal.incrementAndGet();
log.error("직렬화 실패 — 재시도 불가. orderId={}", orderId, e);
quarantine(orderId, payload, e);
} catch (KafkaException e) {
// 버퍼 고갈(max.block.ms 만료)·메타데이터 미수신 등이 여기로 옵니다.
// 호출자에게 알려 백프레셔를 걸 수 있게 다시 던집니다.
failed.incrementAndGet();
log.error("send() 호출 실패 — 상류에 백프레셔가 필요합니다. orderId={}", orderId, e);
throw e;
}
}
/**
* 콜백에서 온 실패를 성격별로 나눕니다.
*
* 여기까지 온 retriable 예외는 이미 delivery.timeout.ms 안에서
* 재시도를 모두 소진한 것입니다. 즉 "재시도 가능"이라는 이름과 달리
* 애플리케이션이 다시 시도해야 하는 상태입니다.
*/
private void handleSendFailure(String orderId, String payload, Exception exception) {
if (exception instanceof RetriableException) {
// 예: NotEnoughReplicasException(min.insync.replicas 미달),
// TimeoutException, NotLeaderOrFollowerException
// 클라이언트 재시도 예산을 다 쓴 뒤이므로 유실로 취급하고 반드시 보고합니다.
failed.incrementAndGet();
log.error("재시도 예산 소진 — 유실 위험. orderId={} cause={}",
orderId, exception.getClass().getSimpleName(), exception);
quarantine(orderId, payload, exception);
return;
}
if (exception instanceof RecordTooLargeException) {
// 브로커 message.max.bytes(기본 1048588, 압축 후 배치 기준) 또는
// 프로듀서 max.request.size(기본 1048576) 초과.
// 같은 레코드를 다시 보내도 항상 실패합니다.
fatal.incrementAndGet();
log.error("레코드가 너무 큽니다 — 재시도 무의미. orderId={}", orderId, exception);
quarantine(orderId, payload, exception);
return;
}
if (exception instanceof AuthorizationException
|| exception instanceof UnknownTopicOrPartitionException) {
// 설정·권한 문제입니다. 재시도로 해결되지 않으므로
// 즉시 크게 알려서 사람이 개입하게 만듭니다.
fatal.incrementAndGet();
log.error("설정/권한 오류 — 즉시 조치 필요. orderId={}", orderId, exception);
quarantine(orderId, payload, exception);
return;
}
// 분류되지 않은 예외도 절대 무시하지 않습니다.
failed.incrementAndGet();
log.error("분류되지 않은 전송 실패. orderId={}", orderId, exception);
quarantine(orderId, payload, exception);
}
/**
* 실패한 이벤트를 잃지 않도록 격리합니다.
*
* 실제 서비스에서는 로컬 디스크 큐 · 관계형 DB 의 outbox 테이블 ·
* 별도 재처리 토픽 중 하나를 씁니다.
* 중요한 것은 "여기서 아무것도 하지 않으면 그 이벤트는 사라진다" 는 점입니다.
* 이 예제에서는 표준 에러로 남기는 것으로 대체합니다.
*/
private void quarantine(String orderId, String payload, Exception cause) {
log.error("QUARANTINE orderId={} payload={} reason={}",
orderId, payload, cause.getClass().getName());
}
/** 지금까지의 집계. 모니터링·테스트에서 씁니다. */
public String stats() {
return "sent=%d, failed=%d, fatal=%d".formatted(sent.get(), failed.get(), fatal.get());
}
public long failedCount() {
return failed.get() + fatal.get();
}
/**
* 종료 처리.
*
* flush() 는 버퍼에 남은 모든 레코드가 성공 또는 실패로 확정될 때까지 블록합니다.
* 이 호출을 빼먹으면 프로세스가 끝날 때 버퍼의 레코드가 그대로 사라집니다.
* close(Duration) 도 내부적으로 flush 하지만, 타임아웃이 지나면 남은 것을 버립니다.
* 그래서 flush() 를 먼저 호출해 "다 보냈다" 를 보장한 뒤 닫습니다.
*/
@Override
public void close() {
try {
producer.flush();
} catch (KafkaException e) {
log.error("flush 실패 — 버퍼에 남은 레코드가 유실될 수 있습니다", e);
} finally {
// 타임아웃을 넉넉히 둡니다. delivery.timeout.ms(120초)보다 짧으면
// 아직 재시도 중인 레코드를 잘라 버립니다.
producer.close(Duration.ofSeconds(130));
}
log.info("프로듀서 종료. {}", stats());
}
}
진입점
package com.example.kafka;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.UUID;
public final class Main {
private static final Logger log = LoggerFactory.getLogger(Main.class);
// 예제 1의 클러스터. 호스트에서 접속하는 PLAINTEXT_HOST 리스너 주소입니다.
private static final String BOOTSTRAP =
"localhost:29092,localhost:39092,localhost:49092";
private static final String TOPIC = "orders";
public static void main(String[] args) {
int count = args.length > 0 ? Integer.parseInt(args[0]) : 1_000;
// 1) 토픽이 요구 스펙으로 존재하는지 먼저 보장합니다.
TopicProvisioner.ensureTopic(BOOTSTRAP, TOPIC, 6, (short) 3, "2");
// 2) try-with-resources 로 종료 시 flush + close 를 보장합니다.
try (OrderApprovedProducer producer =
new OrderApprovedProducer(BOOTSTRAP, TOPIC, "order-approved-producer")) {
// JVM 이 SIGTERM 을 받아도 버퍼를 비우도록 훅을 등록합니다.
// 컨테이너 환경에서는 이것이 없으면 배포마다 소량 유실이 발생합니다.
Runtime.getRuntime().addShutdownHook(new Thread(producer::close, "producer-shutdown"));
for (int i = 0; i < count; i++) {
String orderId = "ORD-" + UUID.randomUUID();
String payload = """
{"orderId":"%s","amount":%d,"currency":"KRW","status":"APPROVED"}"""
.formatted(orderId, 10_000 + i);
producer.send(orderId, payload);
}
log.info("전송 요청 완료: {}건. flush 를 기다립니다.", count);
} // ← 여기서 flush() → close() 가 실행됩니다.
log.info("정상 종료");
}
}
실행 방법
# 0. 예제 1의 클러스터가 떠 있어야 합니다.
cd kafka-lab && docker compose ps
cd ..
# 1. 빌드
cd safe-producer
mvn -q clean package
# 2. 1000건 전송
mvn -q exec:java -Dexec.args=1000
# 3. (선택) 직접 java 로 실행하려면 의존성을 한 디렉터리에 모읍니다.
mvn -q dependency:copy-dependencies -DoutputDirectory=target/lib
java -cp "target/safe-producer-1.0.0.jar:target/lib/*" com.example.kafka.Main 1000
[main] INFO com.example.kafka.TopicProvisioner - 토픽이 이미 존재합니다: orders — 스펙을 검증합니다
[main] INFO org.apache.kafka.clients.producer.ProducerConfig - ProducerConfig values:
acks = all
batch.size = 65536
compression.type = lz4
delivery.timeout.ms = 120000
enable.idempotence = true
linger.ms = 20
max.in.flight.requests.per.connection = 5
retries = 2147483647
[kafka-producer-network-thread | order-approved-producer] INFO org.apache.kafka.clients.producer.internals.TransactionManager - [Producer clientId=order-approved-producer] ProducerId set to 4000 with epoch 0
[main] INFO com.example.kafka.Main - 전송 요청 완료: 1000건. flush 를 기다립니다.
[main] INFO com.example.kafka.OrderApprovedProducer - 프로듀서 종료. sent=1000, failed=0, fatal=0
검증 방법
1. 건수가 정확히 맞는가
cd kafka-lab
# 각 파티션의 log-end-offset 을 조회합니다(-1 은 latest 를 의미).
./kcli kafka-get-offsets.sh --topic orders --time -1
orders:0:169
orders:1:158
orders:2:171
orders:3:165
orders:4:172
orders:5:165
키가 ORD-{UUID}로 매번 다르므로 파티션에 고르게 퍼집니다.
합계 1000이 나오면 브로커가 실제로 1000건을 커밋했다는 뜻입니다.
애플리케이션 로그의 sent=1000과 이 합계가 일치하는지가 핵심입니다.
둘 중 하나만 맞으면 어딘가에서 실패를 삼킨 것입니다.
2. 브로커 1대를 죽여도 성공하는가
# 브로커 1대 정지 → ISR 이 2로 줄어듭니다.
docker compose stop kafka-3
./kcli kafka-topics.sh --describe --topic orders --under-replicated-partitions
# min.insync.replicas=2 이므로 acks=all 쓰기는 여전히 성공해야 합니다.
cd ../safe-producer && mvn -q exec:java -Dexec.args=200
# → sent=200, failed=0
3. 브로커 2대를 죽이면 실패로 보고하는가
이것이 이 예제의 진짜 검증 항목입니다. 유실이 아니라 실패가 나와야 합니다.
cd ../kafka-lab
docker compose stop kafka-2 # 이제 살아 있는 브로커는 kafka-1 뿐입니다
cd ../safe-producer && mvn -q exec:java -Dexec.args=10
[kafka-producer-network-thread | order-approved-producer] WARN o.a.k.c.p.internals.Sender - [Producer clientId=order-approved-producer] Got error produce response with correlation id 8 on topic-partition orders-2, retrying (2147483646 attempts left). Error: NOT_ENOUGH_REPLICAS
...
[kafka-producer-network-thread | order-approved-producer] ERROR com.example.kafka.OrderApprovedProducer - 재시도 예산 소진 — 유실 위험. orderId=ORD-... cause=TimeoutException
[kafka-producer-network-thread | order-approved-producer] ERROR com.example.kafka.OrderApprovedProducer - QUARANTINE orderId=ORD-... payload={...} reason=org.apache.kafka.common.errors.TimeoutException
[main] INFO com.example.kafka.OrderApprovedProducer - 프로듀서 종료. sent=0, failed=10, fatal=0
브로커는 NOT_ENOUGH_REPLICAS 에러를 반환하고,
클라이언트는 이를 retriable로 판단해 재시도합니다.
delivery.timeout.ms=120000이 만료되면 콜백에 예외가 전달되고,
우리 코드가 QUARANTINE으로 격리합니다.
같은 상황에서 acks=1이었다면 이 10건은 성공으로 보고되고
kafka-1이 죽는 순간 사라집니다.
cd ../kafka-lab && docker compose start kafka-2 kafka-3
./kcli kafka-topics.sh --describe --topic orders # Isr 이 3으로 회복되는지 확인
4. 멱등성이 중복을 막는가
재시도로 인한 중복이 제거되는지는 중복 카운트로 확인합니다.
위의 3번 실험을 min.insync.replicas=1로 바꿔 반복하면
재시도가 발생하면서도 최종 건수가 정확히 맞는 것을 볼 수 있습니다.
# 키만 뽑아서 중복이 있는지 셉니다. 출력이 없으면 중복이 없다는 뜻입니다.
./kcli kafka-console-consumer.sh --topic orders --from-beginning \
--timeout-ms 20000 \
--property print.key=true --property print.value=false \
2>/dev/null | sort | uniq -d
이 예제에서 쓴 설정
| 설정 | 기본값 | 이 예제 값 | 이유 / 튜닝 포인트 |
|---|---|---|---|
acksproducer |
all |
all |
ISR 전부의 기록을 기다립니다. 3.0부터 기본값이 1→all로 바뀌었습니다 |
enable.idempotenceproducer |
true |
true |
3.0부터 기본값. 켜면 acks=all·retries=MAX·in-flight≤5가 강제됩니다 |
min.insync.replicastopic |
1 |
2 |
RF=3에 2가 표준. 3으로 올리면 1대 정비만으로 쓰기가 멈춥니다(케이스 6) |
max.in.flight.requests.per.connectionproducer |
5 |
5 |
멱등성이 켜진 상태에서 6 이상이면 ConfigException으로 기동 실패합니다 |
retriesproducer |
2147483647 |
지정하지 않음 | 4.x에서는 이 값을 만지지 않는 것이 권장입니다. 실질 상한은 delivery.timeout.ms입니다 |
delivery.timeout.msproducer |
120000 |
120000 |
send()부터 성공/실패 확정까지의 총 예산. linger.ms + request.timeout.ms 이상이어야 합니다 |
request.timeout.msproducer |
30000 |
30000 |
요청 1회의 응답 대기. 이 단위로 재시도가 반복됩니다 |
max.block.msproducer |
60000 |
60000 |
버퍼 고갈·메타데이터 미수신 시 send()가 블록되는 상한 |
linger.msproducer |
5 |
20 |
4.0에서 0→5로 변경. 시중 자료 대부분이 0으로 적혀 있습니다 |
batch.sizeproducer |
16384 |
65536 |
바이트 단위. 파티션마다 이 크기의 버퍼가 생기므로 무작정 키우면 메모리를 먹습니다 |
buffer.memoryproducer |
33554432 |
67108864 |
전송 대기 총량. 고갈 시 send()가 max.block.ms만큼 블록됩니다 |
compression.typeproducer |
none |
lz4 |
브로커의 message.max.bytes(1048588)는 압축 후 배치 크기 기준입니다 |
max.request.sizeproducer |
1048576 |
지정하지 않음 | 브로커 message.max.bytes(1048588)와 값이 미묘하게 다릅니다(케이스 10) |
unclean.leader.election.enabletopic |
false |
기본값 유지 | true면 ISR 밖 레플리카가 리더가 되어 데이터를 잘라냅니다. 절대 켜지 마세요(케이스 3) |
delivery.timeout.ms 예산 —
request.timeout.ms×재시도가 delivery.timeout.ms 안에 들어가는 포함 관계
프로덕션 고려사항
| 항목 | 이 예제 | 프로덕션 |
|---|---|---|
| 격리(quarantine) 구현 | 로그로 남기기만 합니다 | outbox 테이블 또는 로컬 디스크 큐에 원자적으로 저장하고 별도 워커가 재발행합니다. 로그만 남기면 사람이 못 보고 지나갑니다 |
| 이벤트 생성과 발행의 원자성 | DB 없음 | DB 커밋과 Kafka 발행은 서로 다른 트랜잭션입니다. 둘을 묶으려면 outbox 패턴 + CDC(예제 8)를 씁니다 |
| 직렬화 | StringSerializer + 손으로 만든 JSON |
Avro/Protobuf + Schema Registry. 스키마 없는 JSON은 컨슈머 배포 순서 사고를 만듭니다(예제 7) |
| 모니터링 | 애플리케이션 카운터 | record-error-rate·record-retry-rate·buffer-available-bytes를 JMX로 수집하고 record-error-rate > 0에 알림(예제 10) |
| 프로듀서 인스턴스 수 | 1개 | KafkaProducer는 스레드 세이프합니다. 스레드마다 만들지 말고 하나를 공유하는 편이 배치 효율이 좋습니다 |
| 종료 처리 | shutdown hook + flush() |
K8s라면 terminationGracePeriodSeconds를 delivery.timeout.ms보다 크게 잡아야 합니다. 그러지 않으면 SIGKILL로 버퍼가 날아갑니다 |
| 백프레셔 | 예외를 던지고 끝 | 버퍼 고갈 시 상류(HTTP 요청 등)에 429를 반환하거나 큐 길이를 제한합니다. 무한정 send()하면 OOM으로 갑니다 |
자주 하는 실수
관련 케이스 스터디
이어서 볼 곳
공식 문서 출처
- Producer Configs —
acks=all,enable.idempotence=true,linger.ms=5,batch.size=16384,retries=2147483647,delivery.timeout.ms=120000,request.timeout.ms=30000,max.block.ms=60000,buffer.memory=33554432,max.request.size=1048576,compression.type=none,max.in.flight.requests.per.connection=5 - Topic Configs —
min.insync.replicas=1,unclean.leader.election.enable=false,max.message.bytes=1048588 - Broker Configs —
message.max.bytes=1048588 - KafkaProducer Javadoc —
send()의 비동기 의미, 멱등성 활성화 시retries가Integer.MAX_VALUE가 된다는 서술, "멱등성을 쓸 때는 애플리케이션 레벨 재전송을 피하라"는 권고, 트랜잭션 토픽에 RF≥3·min.insync.replicas=2권고 - Message Delivery Semantics — at-least-once / exactly-once의 정의와 조건