학습 목표

시나리오

주문 이벤트를 읽어 분석용 데이터 웨어하우스에 배치 적재하는 컨슈머입니다. 건당 INSERT는 비효율적이라 500건씩 모아 한 번에 COPY합니다. orders 토픽은 파티션 6개, 컨슈머 인스턴스 3대입니다.

운영 중 두 번의 사고가 있었습니다.

요구사항은 "유실은 절대 금지, 중복은 최소화하고 적재 쪽에서 멱등 처리"입니다.

아키텍처

커밋 시점별 유실과 중복 — 처리 전 커밋은 유실, 처리 후 커밋은 중복 poll 로 오프셋 5 부터 9 까지 다섯 건을 받은 컨슈머가 도중에 크래시하는 상황을 두 가지 커밋 순서로 비교합니다. 첫째는 처리 전에 커밋하는 경우입니다. 오프셋 10 을 먼저 커밋한 뒤 5 와 6 만 처리하고 크래시하면, 재시작한 컨슈머는 커밋된 10 부터 읽으므로 7, 8, 9 는 아무도 처리하지 않은 채 건너뛰게 됩니다. 이것이 유실이며 at-most-once 성질입니다. 둘째는 처리 후에 커밋하는 경우입니다. 5 부터 9 까지 처리를 마쳤지만 커밋 직전에 크래시하면, 재시작한 컨슈머는 이전 커밋 지점인 5 부터 다시 읽으므로 다섯 건을 다시 처리합니다. 이것이 중복이며 at-least-once 성질입니다. 기본값 enable.auto.commit 이 true 이면 poll 을 부를 때 주기적으로 커밋되므로, 처리가 끝나기 전에 커밋이 나가 첫째 경우가 될 수 있습니다. 중복을 없애려면 커밋 순서만으로는 부족하고, 컨슈머 쪽 처리를 멱등하게 만들거나 처리 결과와 오프셋을 한 트랜잭션으로 묶어야 합니다. 같은 크래시, 커밋 순서만 다릅니다 — 결과는 유실과 중복으로 갈립니다 ① 처리 커밋 — enable.auto.commit=true(기본) 에서 흔히 생기는 모양 poll() offset 5 ~ 9 5건 받음 커밋 committed = 10 아직 처리 안 함 처리 중 크래시 5, 6 만 처리됨 7, 8, 9 미처리 재시작 10 부터 읽음 건너뜀 결과 5 ✔ 6 ✔ 7 ✕ 8 ✕ 9 ✕ 3건 유실 · 아무도 처리하지 않았습니다 (at-most-once) ② 처리 커밋 — enable.auto.commit=false + 처리 완료 후 commitSync() poll() offset 5 ~ 9 이전 커밋 = 5 처리 완료 5 ~ 9 모두 처리 부수 효과 발생 커밋 직전 크래시 committed 은 5 커밋이 안 나감 재시작 5 부터 읽음 재처리 결과 5 ×2 6 ×2 7 ×2 8 ×2 9 ×2 5건 중복 처리 · 유실은 없습니다 (at-least-once) 유실과 중복 중 하나는 반드시 고릅니다. 대부분의 업무는 ②(중복 허용)를 고르고 처리를 멱등하게 만듭니다. 멱등하게 만드는 방법: 레코드 키를 업무 키로 쓰기 · upsert 사용 · 처리한 오프셋을 결과와 같은 트랜잭션에 기록 enable.auto.commit=truepoll() 호출 시점에 auto.commit.interval.ms=5000 주기로 커밋합니다. 즉 자동 커밋은 ①과 ② 사이 어디든 될 수 있습니다. 유실을 막아야 하면 자동 커밋을 끄고 직접 커밋하세요.
커밋 시점별 유실·중복 — 처리 전 커밋(유실)과 처리 후 커밋(중복)의 차이
poll 루프 내부 동작 — poll() 은 대개 네트워크 호출이 아니라 버퍼에서 꺼내는 일 컨슈머 내부를 두 구역으로 나눈 그림입니다. 위 구역은 데이터를 미리 받아 두는 경로입니다. 컨슈머 내부의 Fetcher 가 브로커에 FetchRequest 를 보내면 브로커는 fetch.min.bytes 만큼 데이터가 모이거나 fetch.max.wait.ms 가 지날 때까지 기다린 뒤 응답하고, 응답으로 온 레코드 배치들이 컨슈머 안의 완료된 fetch 큐에 쌓입니다. 아래 구역은 poll 이 무엇을 반환하는지 두 경우로 보여 줍니다. 첫째 경우는 큐에 이미 레코드가 있는 경우로, poll 은 네트워크 왕복 없이 큐에서 최대 max.poll.records 500 건을 꺼내 즉시 반환합니다. 둘째 경우는 큐가 빈 경우로, fetch 응답이 도착할 때까지 기다리고 timeout 안에 오지 않으면 빈 목록을 반환합니다. max.poll.records 는 한 번에 몇 건을 돌려줄지만 정하고 브로커에서 얼마나 받아오는지는 바꾸지 않습니다. 받아온 레코드는 캐시에 두고 여러 번의 poll 에 나눠 반환합니다. 반환된 레코드를 처리하고 오프셋을 커밋한 뒤 다시 poll 을 부르는 것이 컨슈머 루프이며, poll 을 계속 부르는 것 자체가 그룹 멤버십 유지 조건입니다. poll() 이 반환하는 곳은 내부 버퍼 입니다. 브로커까지 가는 것은 그 앞 단계입니다. 구역 A · 데이터를 미리 받아 큐에 쌓는 경로 Fetcher 컨슈머 내부에서 fetch 요청을 보냄 fetch 브로커 fetch.min.bytes=1 만큼 모이거나 fetch.max.wait.ms=500 까지 대기 배치 완료된 fetch 큐 (내부 버퍼) 받아온 레코드가 여기 쌓입니다 한도: max.partition.fetch.bytes 구역 B · poll(Duration) 이 반환하는 두 경우 poll() 경우 ① 큐에 레코드가 있음 네트워크 왕복 없음 최대 max.poll.records=500 건 반환 즉시 돌아옵니다 poll() 경우 ② 큐가 비어 있음 fetch 응답을 기다림 도착분 반환 · 안 오면 빈 목록 poll(Duration) 인자만큼만 기다립니다 반환된 레코드 처리 애플리케이션 스레드 오프셋 커밋 자동 또는 수동 다시 poll() 늦으면 그룹에서 빠집니다 max.poll.records 는 fetch 크기를 바꾸지 않습니다. 받아 둔 레코드를 몇 건씩 나눠 줄지만 정합니다. 그래서 이 값을 줄여도 네트워크 사용량은 그대로이고, 한 번의 처리 시간만 짧아집니다. 받아오는 양은 fetch.max.bytes=52428800max.partition.fetch.bytes=1048576 가 정합니다.
poll 루프 내부 동작 — fetch 요청 → 내부 큐 → max.poll.records만큼 반환 → 처리 → 커밋
커밋 전략별 보장
전략 보장 최악의 경우 영향 범위 적합한 워크로드
자동 커밋
(enable.auto.commit=true)
사실상 at-most-once에 가까움 auto.commit.interval.ms(기본 5000) 동안 처리한 것이 유실될 수 있습니다 유실이 허용되는 메트릭·로그 수집
처리 전 수동 커밋 at-most-once 커밋한 배치 전체 유실 중복이 절대 안 되는 경우 (드묾)
처리 후 수동 커밋 이 예제 at-least-once 마지막 배치만큼 중복 (최대 max.poll.records건) 대부분의 경우. 하류에서 멱등 처리
트랜잭션 + sendOffsetsToTransaction exactly-once (Kafka 경계 안) 없음 출력이 Kafka 토픽인 경우 (예제 5)

사전 요구사항

전체 코드

디렉터리 구조
offset-strategy/
├── pom.xml
└── src/main/java/com/example/kafka/consumer/
    ├── ConsumerConfigFactory.java   # 설정 (classic / consumer 프로토콜 둘 다)
    ├── BatchSink.java               # 배치 적재 대상 (멱등)
    ├── OffsetCommitRebalanceListener.java # 리밸런스 대응
    ├── BatchingConsumer.java        # 컨슈머 본체
    └── Main.java
offset-strategy/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>offset-strategy</artifactId>
  <version>1.0.0</version>

  <properties>
    <maven.compiler.release>17</maven.compiler.release>
    <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
    <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>
    <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.consumer.Main</mainClass>
        </configuration>
      </plugin>
    </plugins>
  </build>
</project>

설정 — 두 프로토콜을 한 곳에서

src/main/java/com/example/kafka/consumer/ConsumerConfigFactory.java
package com.example.kafka.consumer;

import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.common.serialization.StringDeserializer;

import java.util.Properties;

public final class ConsumerConfigFactory {

    private ConsumerConfigFactory() {
    }

    /**
     * @param newProtocol true 면 KIP-848 새 컨슈머 그룹 프로토콜(group.protocol=consumer)을 씁니다.
     *                    4.3의 기본값은 여전히 classic 입니다.
     */
    public static Properties create(String bootstrapServers, String groupId,
                                    String clientId, boolean newProtocol) {
        Properties p = new Properties();
        p.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
        p.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
        p.put(ConsumerConfig.CLIENT_ID_CONFIG, clientId);
        p.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
        p.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());

        // --- 커밋 (이 예제의 핵심) -------------------------------------------
        // 자동 커밋을 끕니다. Kafka 기본값은 true 이며,
        // 그대로 두면 auto.commit.interval.ms(기본 5000) 마다
        // "처리 완료와 무관하게" 오프셋이 전진해 유실 창이 생깁니다.
        p.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);

        // 커밋된 오프셋이 없을 때만 적용됩니다.
        // 기본값 latest 는 "새 그룹은 과거 데이터를 건너뛴다" 를 의미합니다.
        // 적재 파이프라인은 과거 데이터도 필요하므로 earliest 로 둡니다.
        // none 으로 두면 커밋된 오프셋이 없을 때 예외로 기동을 막습니다.
        p.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");

        // --- 배치 크기와 시간 예산 -------------------------------------------
        // 한 번의 poll() 이 반환하는 최대 건수. 기본값 500.
        // 이 값이 "최악의 경우 중복되는 건수" 의 상한이기도 합니다.
        p.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 500);

        // poll 사이 최대 간격. 기본값 300000(5분).
        // 배치 처리(웨어하우스 COPY)가 이 시간을 넘으면
        // 브로커가 이 컨슈머를 죽은 것으로 보고 그룹에서 축출합니다.
        // 그러면 다른 인스턴스가 같은 구간을 다시 처리해 중복이 커집니다.
        p.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 300_000);

        // fetch 튜닝: 최소 이만큼 모일 때까지 브로커가 응답을 지연시킵니다.
        // 기본값 1(바이트)이면 사실상 즉시 응답합니다.
        // 배치 적재는 지연에 관대하므로 크게 잡아 요청 수를 줄입니다.
        p.put(ConsumerConfig.FETCH_MIN_BYTES_CONFIG, 64 * 1024);
        // 위 조건이 안 차도 이 시간이면 응답합니다. 기본값 500.
        p.put(ConsumerConfig.FETCH_MAX_WAIT_MS_CONFIG, 500);
        // 파티션당 한 번에 가져오는 최대 바이트. 기본값 1048576(1MB).
        p.put(ConsumerConfig.MAX_PARTITION_FETCH_BYTES_CONFIG, 2 * 1024 * 1024);

        // 상류가 트랜잭션 프로듀서라면 read_committed 가 필요합니다(예제 5).
        // 기본값은 read_uncommitted 입니다.
        p.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_uncommitted");

        if (newProtocol) {
            // --- KIP-848 새 컨슈머 그룹 프로토콜 (4.0 GA, 기본값 아님) --------
            // 브로커가 파티션 할당을 계산하고, 증분 설계라
            // 전역 동기화 장벽(stop-the-world)이 없습니다.
            p.put(ConsumerConfig.GROUP_PROTOCOL_CONFIG, "consumer");

            // 서버 측 할당자를 고를 수 있습니다. 지정하지 않으면
            // 브로커의 group.consumer.assignors 목록 첫 번째(uniform)가 쓰입니다.
            // p.put(ConsumerConfig.GROUP_REMOTE_ASSIGNOR_CONFIG, "uniform");

            // 아래 3개는 새 프로토콜에서 "쓸 수 없습니다".
            // 하트비트와 세션 타임아웃은 브로커의
            // group.consumer.heartbeat.interval.ms / group.consumer.session.timeout.ms
            // 가 대신 관리합니다. 여기에 넣으면 무시되거나 오류가 됩니다.
            //   heartbeat.interval.ms
            //   session.timeout.ms
            //   partition.assignment.strategy
            // 그리고 enforceRebalance() API 도 쓸 수 없습니다.
        } else {
            // --- classic 프로토콜 (4.3 기본값) --------------------------------
            // 하트비트가 이 시간 동안 끊기면 그룹에서 축출됩니다. 기본값 45000.
            p.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 45_000);
            // 하트비트 전송 간격. session.timeout.ms 의 1/3 이하가 권장입니다. 기본값 3000.
            p.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, 3_000);
            // 기본값은 [RangeAssignor, CooperativeStickyAssignor] 입니다.
            // CooperativeSticky 만 쓰면 증분 리밸런스가 되어
            // 리밸런스 때 모든 파티션을 회수하지 않습니다(stop-the-world 완화).
            p.put(ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG,
                    "org.apache.kafka.clients.consumer.CooperativeStickyAssignor");
        }

        return p;
    }
}

적재 대상 — 멱등이어야 합니다

src/main/java/com/example/kafka/consumer/BatchSink.java
package com.example.kafka.consumer;

import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.util.List;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;

/**
 * 배치 적재 대상 (실습용).
 *
 * at-least-once 컨슈머는 같은 레코드를 두 번 줄 수 있습니다.
 * 그래서 적재 쪽이 멱등해야 합니다.
 * 여기서는 "토픽-파티션-오프셋" 을 자연 키로 써서 중복을 걸러 냅니다.
 * 실제 웨어하우스라면 그 조합에 UNIQUE 제약을 걸거나 MERGE 를 씁니다.
 */
public class BatchSink {

    private static final Logger log = LoggerFactory.getLogger(BatchSink.class);

    /** 이미 적재한 자연 키. 실제로는 DB 의 UNIQUE 제약이 이 역할을 합니다. */
    private final Set<String> seen = ConcurrentHashMap.newKeySet();

    private long insertedRows = 0;
    private long duplicateRows = 0;
    private long flushes = 0;

    /**
     * 배치를 적재합니다. 예외를 던지면 호출자가 커밋하지 않습니다.
     *
     * @param batch 적재할 레코드 목록
     */
    public void flush(List<ConsumerRecord<String, String>> batch) {
        if (batch.isEmpty()) {
            return;
        }

        // 실제로는 여기서 트랜잭션을 열고 COPY / batch INSERT 를 수행합니다.
        // 부분 성공이 생기지 않도록 반드시 하나의 DB 트랜잭션으로 묶어야 합니다.
        // 부분 성공 후 커밋하면 나머지가 유실됩니다.
        for (ConsumerRecord<String, String> record : batch) {
            String naturalKey = record.topic() + "-" + record.partition() + "@" + record.offset();
            if (seen.add(naturalKey)) {
                insertedRows++;
            } else {
                // 리밸런스나 재시작으로 다시 받은 레코드입니다.
                duplicateRows++;
            }
        }

        flushes++;
        log.info("적재 완료 batch={}건 (누적 inserted={}, duplicate={}, flushes={})",
                batch.size(), insertedRows, duplicateRows, flushes);
    }

    public String stats() {
        return "inserted=%d, duplicate=%d, flushes=%d"
                .formatted(insertedRows, duplicateRows, flushes);
    }

    public long duplicateRows() {
        return duplicateRows;
    }
}

리밸런스 대응 — 중복을 최소화하는 지점

src/main/java/com/example/kafka/consumer/OffsetCommitRebalanceListener.java
package com.example.kafka.consumer;

import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerRebalanceListener;
import org.apache.kafka.clients.consumer.OffsetAndMetadata;
import org.apache.kafka.common.TopicPartition;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.util.Collection;
import java.util.Map;
import java.util.function.Supplier;

/**
 * 파티션 회수 직전에 진행 중인 배치를 적재하고 오프셋을 커밋합니다.
 *
 * 이것이 없으면 리밸런스마다 최대 max.poll.records 건이 중복 처리됩니다.
 * (다른 인스턴스가 마지막 커밋 지점부터 다시 읽기 때문입니다)
 */
public class OffsetCommitRebalanceListener implements ConsumerRebalanceListener {

    private static final Logger log = LoggerFactory.getLogger(OffsetCommitRebalanceListener.class);

    private final Consumer<String, String> consumer;
    /** 진행 중인 배치를 적재하고 커밋할 오프셋을 반환하는 콜백 */
    private final Supplier<Map<TopicPartition, OffsetAndMetadata>> flushAndGetOffsets;

    public OffsetCommitRebalanceListener(
            Consumer<String, String> consumer,
            Supplier<Map<TopicPartition, OffsetAndMetadata>> flushAndGetOffsets) {
        this.consumer = consumer;
        this.flushAndGetOffsets = flushAndGetOffsets;
    }

    /**
     * 파티션을 잃기 직전에 호출됩니다.
     *
     * cooperative 리밸런스에서는 "실제로 잃는 파티션만" 넘어옵니다.
     * eager 리밸런스에서는 전체가 넘어옵니다.
     * 여기서 반드시 동기 커밋(commitSync)을 써야 합니다 —
     * 비동기 커밋은 파티션을 잃은 뒤에 응답이 와서 실패합니다.
     */
    @Override
    public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
        log.info("파티션 회수 예정: {} — 진행 중인 배치를 적재하고 커밋합니다", partitions);
        try {
            Map<TopicPartition, OffsetAndMetadata> offsets = flushAndGetOffsets.get();
            if (!offsets.isEmpty()) {
                consumer.commitSync(offsets);
                log.info("회수 전 커밋 완료: {}", offsets);
            }
        } catch (RuntimeException e) {
            // 커밋 실패는 중복 처리를 의미할 뿐 유실은 아닙니다.
            // 예외를 던지면 리밸런스 자체가 실패하므로 로그만 남깁니다.
            log.warn("회수 전 커밋 실패 — 중복 처리가 발생할 수 있습니다", e);
        }
    }

    /** 새 파티션을 받은 직후 호출됩니다. */
    @Override
    public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
        log.info("파티션 할당: {}", partitions);
        for (TopicPartition tp : partitions) {
            // committed(...) 로 시작 위치를 확인해 두면 조사에 도움이 됩니다.
            // seek() 로 임의 위치부터 읽을 수도 있습니다(재처리 시).
            log.info("  {} 시작 오프셋 = {}", tp, consumer.position(tp));
        }
    }

    /**
     * 파티션을 비정상적으로 잃었을 때 호출됩니다
     * (예: max.poll.interval.ms 초과로 그룹에서 축출).
     *
     * 이 시점에는 이미 파티션 소유권이 없으므로 커밋하면 실패합니다.
     * 커밋을 시도하지 않는 것이 정답입니다.
     */
    @Override
    public void onPartitionsLost(Collection<TopicPartition> partitions) {
        log.warn("파티션 상실(비정상): {} — 커밋하지 않습니다. "
                + "max.poll.interval.ms 초과를 의심하세요", partitions);
    }
}

컨슈머 본체

src/main/java/com/example/kafka/consumer/BatchingConsumer.java
package com.example.kafka.consumer;

import org.apache.kafka.clients.consumer.CloseOptions;
import org.apache.kafka.clients.consumer.CommitFailedException;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.clients.consumer.OffsetAndMetadata;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.errors.WakeupException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.time.Duration;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.atomic.AtomicBoolean;

/**
 * 배치 적재 컨슈머 (at-least-once).
 *
 * 처리 순서:
 *   1) poll() 로 레코드를 받아 버퍼에 모읍니다.
 *   2) 버퍼가 batchSize 에 도달하거나 flushIntervalMs 가 지나면 적재합니다.
 *   3) 적재가 "성공한 뒤에만" 오프셋을 커밋합니다. → 유실 없음.
 *   4) 커밋 전에 죽으면 그 배치를 다시 받습니다. → 중복 가능(적재 쪽 멱등으로 흡수).
 */
public class BatchingConsumer implements Runnable {

    private static final Logger log = LoggerFactory.getLogger(BatchingConsumer.class);
    private static final Duration POLL_TIMEOUT = Duration.ofMillis(500);

    private final org.apache.kafka.clients.consumer.Consumer<String, String> consumer;
    private final BatchSink sink;
    private final String topic;
    private final int batchSize;
    private final long flushIntervalMs;

    private final AtomicBoolean running = new AtomicBoolean(true);

    /** 아직 적재하지 않은 레코드 */
    private final List<ConsumerRecord<String, String>> buffer = new ArrayList<>();
    /** 버퍼에 담긴 레코드로 계산한 "다음에 읽을 오프셋" */
    private final Map<TopicPartition, OffsetAndMetadata> pendingOffsets = new HashMap<>();

    private long lastFlushAt = System.currentTimeMillis();

    public BatchingConsumer(java.util.Properties props, BatchSink sink,
                            String topic, int batchSize, long flushIntervalMs) {
        this.consumer = new KafkaConsumer<>(props);
        this.sink = sink;
        this.topic = topic;
        this.batchSize = batchSize;
        this.flushIntervalMs = flushIntervalMs;
    }

    @Override
    public void run() {
        // 리밸런스 리스너를 함께 등록합니다.
        // 파티션 회수 직전에 flushAndCollectOffsets() 를 호출하도록 연결합니다.
        consumer.subscribe(List.of(topic),
                new OffsetCommitRebalanceListener(consumer, this::flushAndCollectOffsets));

        try {
            while (running.get()) {
                ConsumerRecords<String, String> records = consumer.poll(POLL_TIMEOUT);

                for (ConsumerRecord<String, String> record : records) {
                    buffer.add(record);
                }
                // nextOffsets() 는 파티션별 "다음에 읽을 오프셋"(+ 리더 에포크)을 줍니다.
                // 직접 offset+1 을 계산하면 리더 에포크가 빠집니다.
                pendingOffsets.putAll(records.nextOffsets());

                if (shouldFlush()) {
                    flushAndCommit();
                }
            }
        } catch (WakeupException e) {
            // close() 에서 wakeup() 을 부른 정상 종료 경로입니다.
            log.info("종료 신호 수신");
        } finally {
            // 종료 직전에 남은 버퍼를 적재하고 커밋합니다.
            // 이 처리가 없으면 배포마다 최대 batchSize 건이 중복 처리됩니다.
            try {
                flushAndCommit();
            } catch (RuntimeException e) {
                log.error("종료 시 적재/커밋 실패 — 재시작 후 중복 처리됩니다", e);
            } finally {
                // 4.3에서 close(Duration) 은 deprecated 입니다. CloseOptions 를 씁니다.
                consumer.close(CloseOptions.timeout(Duration.ofSeconds(30)));
            }
            log.info("컨슈머 종료. {}", sink.stats());
        }
    }

    private boolean shouldFlush() {
        if (buffer.isEmpty()) {
            return false;
        }
        boolean full = buffer.size() >= batchSize;
        boolean timedOut = System.currentTimeMillis() - lastFlushAt >= flushIntervalMs;
        return full || timedOut;
    }

    /**
     * 적재 → 커밋. 순서를 바꾸면 유실이 됩니다.
     */
    private void flushAndCommit() {
        Map<TopicPartition, OffsetAndMetadata> offsets = flushAndCollectOffsets();
        if (offsets.isEmpty()) {
            return;
        }

        try {
            // 동기 커밋. 실패하면 예외가 나므로 놓치지 않습니다.
            // 처리량이 중요하면 commitAsync 를 쓰고 종료 시에만 commitSync 를 씁니다
            // (아래 commitAsyncWithFallback 참고).
            consumer.commitSync(offsets, Duration.ofSeconds(15));
            log.debug("커밋 완료: {}", offsets);

        } catch (CommitFailedException e) {
            // 리밸런스로 파티션 소유권을 잃은 경우입니다.
            // 이미 적재는 끝났으므로 유실은 아닙니다. 중복만 발생합니다.
            log.warn("커밋 실패(리밸런스 추정) — 중복 처리가 발생할 수 있습니다", e);
        }
    }

    /**
     * 버퍼를 적재하고, 커밋할 오프셋을 반환합니다.
     * 적재가 실패하면 예외가 나가고 버퍼와 오프셋은 그대로 남습니다(재시도 가능).
     */
    private Map<TopicPartition, OffsetAndMetadata> flushAndCollectOffsets() {
        if (buffer.isEmpty()) {
            return Map.of();
        }
        // 적재가 성공해야만 아래로 내려갑니다.
        sink.flush(List.copyOf(buffer));

        Map<TopicPartition, OffsetAndMetadata> offsets = Map.copyOf(pendingOffsets);
        buffer.clear();
        pendingOffsets.clear();
        lastFlushAt = System.currentTimeMillis();
        return offsets;
    }

    /**
     * 처리량이 중요한 경로용 대안: 평소에는 비동기, 종료 시에만 동기.
     *
     * commitAsync 는 응답을 기다리지 않아 빠르지만 실패해도 재시도하지 않습니다.
     * 그래서 마지막에 commitSync 로 확실히 마무리하는 조합이 표준 패턴입니다.
     */
    @SuppressWarnings("unused")
    private void commitAsyncWithFallback(Map<TopicPartition, OffsetAndMetadata> offsets) {
        consumer.commitAsync(offsets, (committed, exception) -> {
            if (exception != null) {
                // 비동기 커밋 실패는 보통 무시합니다 —
                // 다음 커밋이 더 큰 오프셋을 덮어쓰기 때문입니다.
                // 다만 계속 실패하면 로그로 드러나야 합니다.
                log.warn("비동기 커밋 실패: {}", committed, exception);
            }
        });
    }

    /** 다른 스레드에서 호출합니다. */
    public void shutdown() {
        running.set(false);
        // poll() 에서 블록 중인 컨슈머를 깨웁니다.
        // KafkaConsumer 는 스레드 세이프하지 않지만 wakeup() 은 예외적으로 안전합니다.
        consumer.wakeup();
    }
}
src/main/java/com/example/kafka/consumer/Main.java
package com.example.kafka.consumer;

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

public final class Main {

    private static final Logger log = LoggerFactory.getLogger(Main.class);

    private static final String BOOTSTRAP =
            "localhost:29092,localhost:39092,localhost:49092";
    private static final String TOPIC = "orders";
    private static final String GROUP_ID = "orders-warehouse-loader";

    public static void main(String[] args) throws InterruptedException {
        // 인자로 프로토콜을 고릅니다: classic(기본) 또는 consumer(KIP-848)
        boolean newProtocol = args.length > 0 && "consumer".equals(args[0]);
        String instanceId = args.length > 1 ? args[1] : "0";
        log.info("group.protocol = {}", newProtocol ? "consumer" : "classic");

        BatchSink sink = new BatchSink();
        BatchingConsumer consumer = new BatchingConsumer(
                ConsumerConfigFactory.create(
                        BOOTSTRAP, GROUP_ID, "warehouse-loader-" + instanceId, newProtocol),
                sink,
                TOPIC,
                500,      // batchSize — max.poll.records 와 같게 두었습니다
                5_000L);  // flushIntervalMs — 유입이 적을 때도 5초마다 적재합니다

        Thread worker = new Thread(consumer, "batching-consumer-" + instanceId);

        Runtime.getRuntime().addShutdownHook(new Thread(() -> {
            log.info("SIGTERM 수신 — 남은 배치를 적재하고 커밋합니다");
            consumer.shutdown();
            try {
                // 종료 시 적재 + 커밋이 끝날 시간을 줍니다.
                worker.join(60_000);
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            }
        }, "shutdown-hook"));

        worker.start();
        worker.join();
    }
}

실행 방법

순서대로 실행
# 0. 예제 1의 클러스터
cd kafka-lab && docker compose ps

# 1. 입력 데이터 준비 (예제 3의 프로듀서를 쓰거나 CLI로)
for i in $(seq 1 5000); do
  echo "ORD-$i:{\"orderId\":\"ORD-$i\",\"amount\":$((1000 + i))}"
done | ./kcli kafka-console-producer.sh --topic orders \
        --property parse.key=true --property key.separator=:

# 2. classic 프로토콜로 실행 (4.3 기본값)
cd ../offset-strategy
mvn -q clean package
mvn -q exec:java -Dexec.args="classic 0"

# 3. 다른 터미널에서 두 번째 인스턴스를 띄워 리밸런스를 관찰
mvn -q exec:java -Dexec.args="classic 1"

# 4. KIP-848 새 프로토콜로 실행 (그룹이 비어 있을 때 전환됩니다)
mvn -q exec:java -Dexec.args="consumer 0"

검증 방법

1. 유실이 없는가 — 강제 종료 후 lag 확인

적재 중 SIGKILL
# 처리 중에 강제 종료합니다(종료 훅을 건너뜁니다).
pkill -9 -f 'com.example.kafka.consumer.Main'

# 커밋되지 않은 구간이 lag 으로 남아 있어야 합니다.
cd ../kafka-lab
./kcli kafka-consumer-groups.sh --describe --group orders-warehouse-loader
기대 출력 — LAG이 0이 아니면 재처리 대상이 남아 있다는 뜻입니다
GROUP                     TOPIC   PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG
orders-warehouse-loader   orders  0          500             834             334
orders-warehouse-loader   orders  1          500             821             321

LAG이 남아 있는 것이 정답입니다. 오프셋이 적재 성공 뒤에만 전진하므로 적재되지 않은 구간이 그대로 남습니다. 재기동하면 그 구간부터 다시 읽습니다 — CURRENT-OFFSETLOG-END-OFFSET까지 올라가는지 확인하세요.

2. 자동 커밋과 비교 — 유실이 생기는 것을 확인

ConsumerConfigFactory를 임시로 바꿔 비교
// 임시 변경 (검증용, 절대 프로덕션에 두지 마세요)
p.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, true);
// 기본값 5000. 이 간격이 "유실 창" 의 크기입니다.
p.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, 5_000);

같은 실험(적재 중 SIGKILL)을 반복하면 lag이 0인데 BatchSinkinserted가 그보다 적습니다. 즉 커밋은 되었는데 적재는 안 된 구간이 있습니다. 그 구간은 어떤 에러도 남기지 않고 영구히 사라집니다 — 이것이 시나리오의 첫 번째 사고입니다.

3. 리밸런스 시 중복이 줄어드는가

두 번째 인스턴스를 띄우고 로그 관찰
[batching-consumer-0] INFO  ...OffsetCommitRebalanceListener - 파티션 회수 예정: [orders-3, orders-4, orders-5] — 진행 중인 배치를 적재하고 커밋합니다
[batching-consumer-0] INFO  ...BatchSink - 적재 완료 batch=213건 (누적 inserted=2713, duplicate=0, flushes=7)
[batching-consumer-0] INFO  ...OffsetCommitRebalanceListener - 회수 전 커밋 완료: {orders-3=..., orders-4=..., orders-5=...}
[batching-consumer-1] INFO  ...OffsetCommitRebalanceListener - 파티션 할당: [orders-3, orders-4, orders-5]

duplicate=0이 핵심입니다. onPartitionsRevoked에서 적재+커밋을 하므로 새 인스턴스가 이어받는 지점이 정확합니다. 이 리스너를 제거하고 같은 실험을 반복하면 duplicate가 최대 max.poll.records(500)까지 늘어납니다 — 시나리오의 두 번째 사고입니다.

4. KIP-848 프로토콜이 실제로 적용되었는가

그룹 타입 확인
# 그룹 목록에 타입이 함께 나옵니다.
./kcli kafka-consumer-groups.sh --list

# 상세 — 새 프로토콜 그룹은 Type 이 Consumer 로 표시됩니다.
./kcli kafka-consumer-groups.sh --describe --group orders-warehouse-loader --state

전환 방법은 두 가지입니다.

이 예제에서 쓴 설정

Apache Kafka 4.3 컨슈머 설정
설정 기본값 이 예제 값 이유
enable.auto.committruefalse기본값이 true라는 점이 함정입니다. 처리 완료와 무관하게 커밋이 진행됩니다
auto.commit.interval.ms5000(사용 안 함)자동 커밋을 쓸 때 "유실 창"의 크기입니다
auto.offset.resetlatestearliest커밋된 오프셋이 없을 때만 적용됩니다. 재처리 수단이 아닙니다
max.poll.records500500최악의 경우 중복 건수의 상한이기도 합니다
max.poll.interval.ms300000300000배치 처리 시간이 이 값을 넘으면 그룹에서 축출됩니다
session.timeout.ms4500045000 (classic만)새 프로토콜에서는 쓸 수 없습니다. 브로커의 group.consumer.session.timeout.ms가 대신합니다
heartbeat.interval.ms30003000 (classic만)새 프로토콜에서는 쓸 수 없습니다
partition.assignment.strategy[RangeAssignor,
CooperativeStickyAssignor]
CooperativeStickyAssignor (classic만)새 프로토콜에서는 쓸 수 없습니다. 브로커의 group.consumer.assignors가 결정합니다
group.protocolclassic인자로 선택KIP-848은 4.0 GA지만 기본값이 아닙니다. 시중 자료가 가장 많이 틀리는 지점입니다
fetch.min.bytes165536배치 적재는 지연에 관대하므로 요청 수를 줄입니다
fetch.max.wait.ms500500fetch.min.bytes가 안 차도 이 시간이면 응답합니다
max.partition.fetch.bytes10485762097152파티션당 fetch 상한. 전체 상한은 fetch.max.bytes(52428800)입니다
isolation.levelread_uncommittedread_uncommitted상류가 트랜잭션을 쓰면 read_committed로 바꿔야 합니다

프로덕션 고려사항

로컬 예제와 프로덕션의 차이
항목이 예제프로덕션
적재 원자성 메모리 Set 배치 전체를 하나의 DB 트랜잭션으로 묶어야 합니다. 부분 성공 후 커밋하면 나머지가 유실됩니다
멱등성 자연 키 중복 제거 (topic, partition, offset)에 UNIQUE 제약을 걸거나 MERGE/UPSERT를 씁니다. 애플리케이션 메모리에 의존하면 재시작 시 무력화됩니다
처리 시간 즉시 반환 배치 처리 시간이 max.poll.interval.ms를 넘을 위험이 있으면 pause()/resume()으로 poll을 유지하면서 별도 스레드에서 처리합니다
인스턴스 재시작 매번 리밸런스 group.instance.id(static membership)를 주면 재시작이 session.timeout.ms 안에 끝날 때 리밸런스를 건너뜁니다. 배포가 잦은 서비스에 효과가 큽니다
종료 처리 join(60s) K8s terminationGracePeriodSeconds를 그보다 크게. 짧으면 SIGKILL로 마지막 배치가 중복 처리됩니다
lag 모니터링 CLI records-lag-max 메트릭 + kafka-consumer-groups.sh 결과를 함께 봅니다. 파티션별 lag 편차가 문제를 드러냅니다(예제 10)
프로토콜 전환 인자로 선택 스테이징에서 consumer 프로토콜을 먼저 검증하세요. 커스텀 assignor를 쓰고 있으면 무중단 전환이 불가능합니다

자주 하는 실수

이어서 볼 곳

공식 문서 출처