실무 예제 · 4
컨슈머 오프셋 전략 (수동 커밋 · 배치 처리)
컨슈머의 유실과 중복은 오프셋을 언제 커밋하는지로 결정됩니다. 처리 전에 커밋하면 유실, 처리 후에 커밋하면 중복입니다. 이 예제는 둘 사이의 선택을 코드로 드러내고, 배치 단위 처리·리밸런스 대응·KIP-848 새 프로토콜까지 함께 다룹니다.
학습 목표
- 자동 커밋이 왜 유실을 만드는지,
auto.commit.interval.ms가 그 창(window)의 크기임을 설명할 수 있습니다. commitSync와commitAsync를 조합한 표준 종료 패턴을 작성할 수 있습니다.ConsumerRebalanceListener로 파티션 회수 직전에 커밋해 중복을 줄일 수 있습니다.group.protocol=consumer(KIP-848)로 전환할 때 쓸 수 없게 되는 설정 3개를 알 수 있습니다.
시나리오
주문 이벤트를 읽어 분석용 데이터 웨어하우스에 배치 적재하는 컨슈머입니다.
건당 INSERT는 비효율적이라 500건씩 모아 한 번에 COPY합니다.
orders 토픽은 파티션 6개, 컨슈머 인스턴스 3대입니다.
운영 중 두 번의 사고가 있었습니다.
- 유실 — 자동 커밋을 쓰던 시절, 배포 중 컨테이너가 SIGKILL되어 이미 커밋된 구간의 데이터가 적재되지 않았습니다. 아무 에러도 남지 않았습니다.
- 대량 중복 — 수동 커밋으로 바꾼 뒤, 리밸런스 때 커밋하지 않아 최대 500건씩 중복 적재되었습니다.
요구사항은 "유실은 절대 금지, 중복은 최소화하고 적재 쪽에서 멱등 처리"입니다.
아키텍처
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
<?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>
설정 — 두 프로토콜을 한 곳에서
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;
}
}
적재 대상 — 멱등이어야 합니다
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;
}
}
리밸런스 대응 — 중복을 최소화하는 지점
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);
}
}
컨슈머 본체
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();
}
}
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 확인
# 처리 중에 강제 종료합니다(종료 훅을 건너뜁니다).
pkill -9 -f 'com.example.kafka.consumer.Main'
# 커밋되지 않은 구간이 lag 으로 남아 있어야 합니다.
cd ../kafka-lab
./kcli kafka-consumer-groups.sh --describe --group orders-warehouse-loader
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-OFFSET이 LOG-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인데 BatchSink의 inserted가 그보다 적습니다.
즉 커밋은 되었는데 적재는 안 된 구간이 있습니다.
그 구간은 어떤 에러도 남기지 않고 영구히 사라집니다 —
이것이 시나리오의 첫 번째 사고입니다.
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
전환 방법은 두 가지입니다.
- 오프라인 — 그룹의 모든 컨슈머를 내리면 그룹이 비고,
group.protocol=consumer로 다시 올리면Consumer타입으로 변환됩니다. 다운타임이 필요합니다. - 온라인(무중단) —
group.protocol=consumer컨슈머를 롤링 배포하면 첫 컨슈머가 합류하는 순간 그룹이Classic→Consumer로 변환되고, 두 프로토콜이 상호 운용됩니다. 단 기존 그룹이 커스텀 메타데이터를 담는 assignor를 쓰지 않는 경우에만 가능합니다.
이 예제에서 쓴 설정
| 설정 | 기본값 | 이 예제 값 | 이유 |
|---|---|---|---|
enable.auto.commit | true | false | 기본값이 true라는 점이 함정입니다. 처리 완료와 무관하게 커밋이 진행됩니다 |
auto.commit.interval.ms | 5000 | (사용 안 함) | 자동 커밋을 쓸 때 "유실 창"의 크기입니다 |
auto.offset.reset | latest | earliest | 커밋된 오프셋이 없을 때만 적용됩니다. 재처리 수단이 아닙니다 |
max.poll.records | 500 | 500 | 최악의 경우 중복 건수의 상한이기도 합니다 |
max.poll.interval.ms | 300000 | 300000 | 배치 처리 시간이 이 값을 넘으면 그룹에서 축출됩니다 |
session.timeout.ms | 45000 | 45000 (classic만) | 새 프로토콜에서는 쓸 수 없습니다. 브로커의 group.consumer.session.timeout.ms가 대신합니다 |
heartbeat.interval.ms | 3000 | 3000 (classic만) | 새 프로토콜에서는 쓸 수 없습니다 |
partition.assignment.strategy | [RangeAssignor, | CooperativeStickyAssignor (classic만) | 새 프로토콜에서는 쓸 수 없습니다. 브로커의 group.consumer.assignors가 결정합니다 |
group.protocol | classic | 인자로 선택 | KIP-848은 4.0 GA지만 기본값이 아닙니다. 시중 자료가 가장 많이 틀리는 지점입니다 |
fetch.min.bytes | 1 | 65536 | 배치 적재는 지연에 관대하므로 요청 수를 줄입니다 |
fetch.max.wait.ms | 500 | 500 | fetch.min.bytes가 안 차도 이 시간이면 응답합니다 |
max.partition.fetch.bytes | 1048576 | 2097152 | 파티션당 fetch 상한. 전체 상한은 fetch.max.bytes(52428800)입니다 |
isolation.level | read_uncommitted | read_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를 쓰고 있으면 무중단 전환이 불가능합니다 |
자주 하는 실수
관련 케이스 스터디
이어서 볼 곳
공식 문서 출처
- Consumer Configs —
enable.auto.commit=true,auto.commit.interval.ms=5000,auto.offset.reset=latest,max.poll.records=500,max.poll.interval.ms=300000,session.timeout.ms=45000,heartbeat.interval.ms=3000,fetch.min.bytes=1,fetch.max.wait.ms=500,fetch.max.bytes=52428800,max.partition.fetch.bytes=1048576,isolation.level=read_uncommitted,group.protocol=classic,partition.assignment.strategy=[RangeAssignor, CooperativeStickyAssignor] - Consumer Rebalance Protocol — 새 프로토콜에서 쓸 수 없는 설정 3개와
enforceRebalance(), 온라인/오프라인 업그레이드 절차, 3.7 EA → 4.0 GA → 5.0 기본값 전환 → 6.0classic제거 로드맵, 브로커 측group.consumer.*설정 - Broker Configs —
group.coordinator.rebalance.protocols=classic,consumer,streams - KafkaConsumer Javadoc — 스레드 세이프하지 않다는 서술,
wakeup()이 유일한 예외,close(Duration)이 deprecated이고close(CloseOptions)를 쓰라는 안내 - ConsumerRecords Javadoc —
nextOffsets()가 다음 오프셋과 리더 에포크를 함께 제공