학습 목표

시나리오

결제 게이트웨이가 카드사 승인 결과를 받으면 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 비교 — 각 설정에서 유실이 생기는 지점 같은 구성(리더 1대와 팔로워 2대)에 대해 acks 값 세 가지를 나란히 놓은 시퀀스 비교입니다. 각 열에는 프로듀서가 전송을 완료로 간주하는 지점을 가로 파선으로 표시했습니다. 이 선이 아래로 내려갈수록 보장이 강해집니다. acks=0 은 소켓 버퍼에 쓴 직후 완료로 보므로 요청이 브로커에 닿지 않아도 프로듀서가 알 수 없고 유실됩니다. offset 은 항상 -1 로 돌아오고 retries 설정도 동작하지 않습니다. acks=1 은 리더가 자기 로그에 기록한 직후 완료로 보므로, 복제가 끝나기 전에 리더가 죽으면 그 레코드는 새 리더에 없어 유실됩니다. acks=all 은 현재 ISR 전원이 응답한 뒤 완료로 보므로 리더가 죽어도 팔로워에 데이터가 있어 유실되지 않습니다. 단 ISR 이 1대로 줄어든 상태에서는 acks=all 도 1대만 확인하므로, min.insync.replicas 를 2 이상으로 두어야 레플리카 한 대 손실을 실제로 견딜 수 있습니다. 가로 파선 = 프로듀서가 전송을 완료로 간주하는 지점. 아래로 갈수록 보장이 강합니다. acks=0 확인 없음 프로듀서 완료 지점 · 소켓 버퍼 여기서 사라져도 모릅니다 리더 — 못 받았을 수도 팔로워 2 팔로워 3 유실 지점: 전송 직후 브로커에 닿지 않아도 성공으로 처리됩니다 retries 무효 · offset -1 acks=1 리더만 확인 프로듀서 ack 리더 — 로컬 로그 기록 완료 지점 · 리더 기록 후 ↓ 복제 전에 리더 다운 미복제 미복제 유실 지점: 복제 이전 ack 을 받은 레코드가 새 리더에는 없습니다 성공으로 보고된 뒤 유실 acks=all ISR 전원 확인 프로듀서 ack 리더 — 로컬 로그 기록 복제 완료 팔로워 2 ✔ 팔로워 3 ✔ 완료 지점 · ISR 전원 응답 결과: 무손실 (조건부) 리더가 죽어도 팔로워에 데이터가 남아 있습니다 단, ISR 이 1대면 위험 acks=all 만으로는 무손실이 아닙니다. ISR 이 1대로 줄면 확인 대상도 1대뿐입니다. min.insync.replicas=2 를 함께 두어야 레플리카 1대 손실을 견딥니다. acks 기본값은 all 입니다.
acks 0 / 1 / all 시퀀스 비교 — 각 설정에서 유실이 발생하는 지점
복제와 ISR — ISR 축소가 min.insync.replicas 에 걸리는 순간 replication.factor 3 인 파티션 orders-0 을 세 단계로 보여 줍니다. 1단계는 리더 broker-1 과 팔로워 broker-2, broker-3 이 모두 따라잡아 ISR 이 3 이고 acks=all 쓰기가 성공합니다. 2단계는 broker-3 이 replica.lag.time.max.ms 30000 밀리초 동안 따라오지 못해 ISR 에서 빠지고 ISR 이 2 로 줄지만 min.insync.replicas 가 2 이므로 아직 성공합니다. 3단계는 broker-2 까지 빠져 ISR 이 1 이 되어 min.insync.replicas 미달이 되고, acks=all 프로듀서는 NOT_ENOUGH_REPLICAS 오류를 받습니다. 이때도 리더는 살아 있으므로 컨슈머 읽기는 계속됩니다. replication.factor 는 3 그대로이고 줄어드는 것은 ISR 집합이라는 점이 핵심입니다. 파티션 orders-0 · replication.factor=3 · min.insync.replicas=2 · acks=all ① 정상 — ISR 3 리더 broker-1 · LEO 120 팔로워 broker-2 · LEO 120 팔로워 broker-3 · LEO 120 ISR = {1, 2, 3} → 크기 3 3 ≥ 2 (충족) 쓰기 성공 ISR 3대 전원이 응답 ② ISR 축소 — 아직 충족 리더 broker-1 · LEO 168 팔로워 broker-2 · LEO 168 팔로워 broker-3 · LEO 88 (지연) replica.lag.time.max.ms 초과 ISR = {1, 2} → 크기 2 2 ≥ 2 (경계, 충족) 쓰기 성공 · 여유 0 한 대만 더 빠지면 중단 ③ 미달 — 쓰기 거부 리더 broker-1 · LEO 168 팔로워 broker-2 · 다운 팔로워 broker-3 · 지연 ISR = {1} → 크기 1 1 < 2 (미달) 쓰기 거부 NOT_ENOUGH_REPLICAS 줄어드는 것은 ISR 집합입니다. replication.factor 는 계속 3 이고, 레플리카가 따라잡으면 ISR 에 다시 들어옵니다. 거부되는 것은 쓰기뿐입니다. 리더가 살아 있으므로 컨슈머는 high watermark 까지 계속 읽을 수 있습니다. ③ 을 피하려면 replication.factor=3 + min.insync.replicas=2 조합으로 레플리카 1대 손실을 허용합니다. 쓰기 성공 조건: ISR 크기 ≥ min.insync.replicas 이고, 그 ISR 전원이 응답. 지연 판정 기준은 30000ms 입니다.
복제와 ISR — 팔로워가 지연되면 ISR이 축소되고 min.insync.replicas에 미달하는 과정
acks × min.insync.replicas 조합에 따른 결과 (RF=3 기준)
조합 브로커 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대 정비만으로 중단)

사전 요구사항

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

safe-producer/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>

무손실 설정 — 값마다 이유가 있습니다

src/main/java/com/example/kafka/SafeProducerConfig.java
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 토픽이 조용히 만들어집니다.

src/main/java/com/example/kafka/TopicProvisioner.java
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);
        }
    }
}

프로듀서 본체 — 실패를 절대 삼키지 않습니다

src/main/java/com/example/kafka/OrderApprovedProducer.java
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());
    }
}

진입점

src/main/java/com/example/kafka/Main.java
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
기대 출력 — 6개 파티션 값의 합이 전송 건수와 같아야 합니다
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대를 죽여도 성공하는가

ISR 축소 상태에서 전송
# 브로커 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대를 죽이면 실패로 보고하는가

이것이 이 예제의 진짜 검증 항목입니다. 유실이 아니라 실패가 나와야 합니다.

min.insync.replicas 미달 유발
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

이 예제에서 쓴 설정

Apache Kafka 4.3 기준 기본값. 소속(프로듀서/토픽/브로커)을 함께 표기했습니다.
설정 기본값 이 예제 값 이유 / 튜닝 포인트
acks
producer
all all ISR 전부의 기록을 기다립니다. 3.0부터 기본값이 1all로 바뀌었습니다
enable.idempotence
producer
true true 3.0부터 기본값. 켜면 acks=all·retries=MAX·in-flight≤5가 강제됩니다
min.insync.replicas
topic
1 2 RF=3에 2가 표준. 3으로 올리면 1대 정비만으로 쓰기가 멈춥니다(케이스 6)
max.in.flight.requests.per.connection
producer
5 5 멱등성이 켜진 상태에서 6 이상이면 ConfigException으로 기동 실패합니다
retries
producer
2147483647 지정하지 않음 4.x에서는 이 값을 만지지 않는 것이 권장입니다. 실질 상한은 delivery.timeout.ms입니다
delivery.timeout.ms
producer
120000 120000 send()부터 성공/실패 확정까지의 총 예산. linger.ms + request.timeout.ms 이상이어야 합니다
request.timeout.ms
producer
30000 30000 요청 1회의 응답 대기. 이 단위로 재시도가 반복됩니다
max.block.ms
producer
60000 60000 버퍼 고갈·메타데이터 미수신 시 send()가 블록되는 상한
linger.ms
producer
5 20 4.0에서 05로 변경. 시중 자료 대부분이 0으로 적혀 있습니다
batch.size
producer
16384 65536 바이트 단위. 파티션마다 이 크기의 버퍼가 생기므로 무작정 키우면 메모리를 먹습니다
buffer.memory
producer
33554432 67108864 전송 대기 총량. 고갈 시 send()max.block.ms만큼 블록됩니다
compression.type
producer
none lz4 브로커의 message.max.bytes(1048588)는 압축 후 배치 크기 기준입니다
max.request.size
producer
1048576 지정하지 않음 브로커 message.max.bytes(1048588)와 값이 미묘하게 다릅니다(케이스 10)
unclean.leader.election.enable
topic
false 기본값 유지 true면 ISR 밖 레플리카가 리더가 되어 데이터를 잘라냅니다. 절대 켜지 마세요(케이스 3)
재시도 예산 — delivery.timeout.ms 가 request.timeout.ms 와 재시도를 모두 포함 delivery.timeout.ms 120000 밀리초를 큰 막대로 그리고, 그 안에 request.timeout.ms 30000 밀리초 단위의 시도들을 중첩해 넣은 그림입니다. 바깥 막대는 send() 호출 이후 성공이나 실패를 보고할 최종 기한이고, 안쪽 칸은 각 전송 시도입니다. 30000 밀리초짜리 시도를 세 번 하면 90000 밀리초를 쓰고 30000 밀리초가 남습니다. 남은 예산 안에 성공하지 못하면 콜백에 TimeoutException 이 전달되며, retries 값이 아무리 커도 여기서 끝납니다. linger.ms 기본값 5 밀리초와 retry.backoff.ms 기본값 100 밀리초는 이 축척에서 선 하나보다 얇아 보이지 않지만 예산에 포함됩니다. 공식 문서의 요구 조건은 delivery.timeout.ms 가 request.timeout.ms 와 linger.ms 의 합보다 크거나 같아야 한다는 것입니다. 재시도 횟수가 아니라 delivery.timeout.ms 예산이 종료 시점을 정합니다 바깥 = 전체 예산 · 안쪽 = 개별 시도 delivery.timeout.ms = 120000send() 이후 성공/실패를 보고할 최종 기한 시도 1 30000 시도 2 (재시도) 30000 시도 3 (재시도) 30000 남은 예산 30000 시도 단위 request 0 30000 60000 90000 120000 linger.ms=5retry.backoff.ms=100 도 예산에 포함되지만 이 축척에서는 선보다 얇습니다. delivery.timeout.ms ≥ request.timeout.ms + linger.ms retries 기본값은 2147483647 입니다. 그래도 무한 재시도가 아닌 이유가 이 예산입니다. 예산을 넘기면 콜백에 TimeoutException 이 전달됩니다. retries 대신 이 값으로 조절하세요. 복구 불가 오류는 예산이 남아도 즉시 실패합니다 (예: RecordTooLargeException). request.timeout.msreplica.lag.time.max.ms(30000) 보다 크게 — 불필요한 재시도 중복을 줄입니다.
재시도와 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라면 terminationGracePeriodSecondsdelivery.timeout.ms보다 크게 잡아야 합니다. 그러지 않으면 SIGKILL로 버퍼가 날아갑니다
백프레셔 예외를 던지고 끝 버퍼 고갈 시 상류(HTTP 요청 등)에 429를 반환하거나 큐 길이를 제한합니다. 무한정 send()하면 OOM으로 갑니다

자주 하는 실수

이어서 볼 곳

공식 문서 출처