실무 예제 · 5
Exactly-Once 파이프라인 (consume-transform-produce)
Kafka의 exactly-once는 "중복이 절대 없다"는 마법이 아닙니다. 출력 레코드와 입력 오프셋 커밋을 하나의 원자적 트랜잭션으로 묶는 것이 전부입니다. 그래서 Kafka 토픽 → 처리 → Kafka 토픽 경계 안에서만 성립합니다. 이 예제는 그 경계 안에서 완전히 동작하는 파이프라인을 만들고, 경계를 넘는 순간 무엇이 깨지는지까지 확인합니다.
학습 목표
initTransactions→beginTransaction→send→sendOffsetsToTransaction→commitTransaction순서를 근거와 함께 쓸 수 있습니다.transactional.id가 왜 재시작 후에도 같아야 하는지, 좀비 프로듀서 차단(fencing)이 어떻게 동작하는지 설명할 수 있습니다.isolation.level=read_committed와 LSO(Last Stable Offset)가 컨슈머에게 어떻게 보이는지 직접 관찰할 수 있습니다.- EOS가 외부 DB·API 호출에는 적용되지 않는다는 점을 코드 수준에서 구분할 수 있습니다.
시나리오
주문 수집 서비스가 orders.raw 토픽에 원본 주문을 넣습니다.
중간 파이프라인이 이 원본을 읽어 통화 정규화와 세금 계산을 하고
orders.enriched 토픽에 씁니다.
하류의 정산 시스템은 orders.enriched만 봅니다.
문제는 이 파이프라인이 배포·크래시·리밸런스로 재시작될 때입니다. at-least-once로 짜면 같은 주문의 보강 결과가 두 번 발행되어 정산 금액이 두 배가 됩니다. at-most-once로 짜면 누락됩니다. 정산 시스템에 멱등 처리를 넣는 것이 정석이지만, 레거시라 손댈 수 없다는 제약이 있습니다.
입력과 출력이 모두 Kafka 토픽이므로 이 구간은 EOS의 적용 대상입니다. Kafka 트랜잭션으로 처리합니다.
아키텍처
initTransactions부터 commitTransaction까지,
트랜잭션 코디네이터와 __transaction_state의 역할
핵심 메커니즘 세 가지
| 부품 | 담당하는 문제 | 없으면 무슨 일이 생기는가 |
|---|---|---|
transactional.id + 에포크 |
재시작·좀비 프로세스 차단(fencing) | 네트워크 분리로 살아남은 옛 인스턴스가 계속 써서 중복이 생깁니다 |
sendOffsetsToTransaction |
출력 쓰기와 입력 오프셋 커밋의 원자성 | 둘 사이에서 죽으면 중복 또는 유실이 됩니다 |
isolation.level=read_committed |
하류 컨슈머가 미완료·중단된 트랜잭션을 보지 않게 함 | abort된 레코드까지 읽어 정산이 틀어집니다 |
사전 요구사항
- 예제 1의 3노드 KRaft 클러스터. 트랜잭션 상태 토픽의 복제 계수가 3, 최소 ISR이 2여야 합니다 — 예제 1의 compose가 그렇게 설정합니다.
- 토픽
orders.raw(6파티션·RF3·min.insync.replicas=2),orders.enriched(동일). 예제 1의create-topics.sh가 만듭니다. - Apache Kafka 4.3.1 클라이언트, Java 17 이상
전체 코드
eos-pipeline/
├── pom.xml
└── src/main/java/com/example/kafka/eos/
├── EosConfig.java # 프로듀서/컨슈머 설정
├── OrderEnricher.java # 순수 변환 로직 (부수효과 없음)
├── EosPipeline.java # 트랜잭션 루프 본체
└── Main.java # 진입점 (인스턴스 ID 로 transactional.id 생성)
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>eos-pipeline</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.eos.Main</mainClass>
</configuration>
</plugin>
</plugins>
</build>
</project>
설정 — 트랜잭션에 필요한 값만 정확히
package com.example.kafka.eos;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.common.IsolationLevel;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.apache.kafka.common.serialization.StringSerializer;
import java.util.Properties;
public final class EosConfig {
private EosConfig() {
}
/**
* 트랜잭션 프로듀서 설정.
*
* @param transactionalId 이 인스턴스의 트랜잭션 ID.
* 재시작해도 같아야 하고, 인스턴스끼리는 달라야 합니다.
*/
public static Properties producer(String bootstrapServers, String transactionalId) {
Properties p = new Properties();
p.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
p.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
p.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
p.put(ProducerConfig.CLIENT_ID_CONFIG, transactionalId);
// --- 트랜잭션 (핵심) -------------------------------------------------
// 이 값을 설정하는 것만으로 트랜잭션 API 가 활성화되고
// enable.idempotence=true, acks=all, retries=MAX 가 자동으로 강제됩니다.
// 재시작 후에도 같은 값을 써야 브로커가 "같은 논리적 프로듀서" 로 인식해
// 이전 세션의 미완료 트랜잭션을 정리(abort)하고 에포크를 올려
// 옛 인스턴스를 차단(fence)할 수 있습니다.
p.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, transactionalId);
// 하나의 트랜잭션이 열려 있을 수 있는 최대 시간. 기본값 60000(1분).
// 이 시간을 넘기면 코디네이터가 트랜잭션을 강제 abort 합니다.
// 브로커의 transaction.max.timeout.ms(기본 900000) 를 넘길 수 없습니다.
// 주의: 한 번의 poll 배치를 처리하는 시간이 이 값을 넘으면
// 커밋이 실패하므로 max.poll.records 와 함께 조율해야 합니다.
p.put(ProducerConfig.TRANSACTION_TIMEOUT_CONFIG, 60_000);
// 내구성 — 트랜잭션이라도 브로커 측 복제가 부실하면 의미가 없습니다.
// transactional.id 를 설정하면 자동으로 all 이 되지만 의도를 남깁니다.
p.put(ProducerConfig.ACKS_CONFIG, "all");
p.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
// 트랜잭션 커밋 자체가 배치 경계이므로 linger 를 크게 둘 이유가 없습니다.
// 4.3 기본값은 5입니다.
p.put(ProducerConfig.LINGER_MS_CONFIG, 5);
p.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "lz4");
return p;
}
/**
* 입력 컨슈머 설정.
*/
public static Properties consumer(String bootstrapServers, String groupId, String clientId) {
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());
// --- EOS 필수 조건 ---------------------------------------------------
// 자동 커밋을 반드시 끕니다. 켜져 있으면 트랜잭션 밖에서 오프셋이 커밋되어
// 원자성이 깨집니다. 공식 문서도 sendOffsetsToTransaction 을 쓸 때는
// 자동 커밋을 끄고 수동 커밋도 하지 말라고 명시합니다.
p.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
// 입력 토픽도 다른 트랜잭션 프로듀서가 쓴다면 read_committed 가 필요합니다.
// 기본값은 read_uncommitted 입니다. 파이프라인을 여러 단으로 잇는다면
// 각 단이 read_committed 여야 합니다.
p.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG,
IsolationLevel.READ_COMMITTED.toString().toLowerCase(java.util.Locale.ROOT));
// 새 그룹일 때 처음부터 읽습니다. 기본값은 latest 입니다.
p.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
// 한 트랜잭션에서 처리할 최대 건수.
// 기본값 500 이지만 트랜잭션 처리 시간이 transaction.timeout.ms 를
// 넘지 않도록 보수적으로 낮춥니다.
p.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 200);
// poll 사이 최대 간격. 기본값 300000(5분).
// 한 배치 처리 시간이 이 값을 넘으면 그룹에서 축출되어 리밸런스가 일어납니다.
p.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 300_000);
return p;
}
}
변환 로직 — 부수효과가 없어야 합니다
트랜잭션 안에서 실행되는 변환은 재실행 가능(idempotent)해야 합니다. 트랜잭션이 abort되면 같은 입력이 다시 처리되기 때문입니다. 여기서 외부 API를 호출하면 그 호출은 롤백되지 않습니다. 그 경계 문제를 아래 EOS의 경계에서 다룹니다.
package com.example.kafka.eos;
import java.math.BigDecimal;
import java.math.RoundingMode;
/**
* 주문 보강 — 순수 함수입니다.
*
* 같은 입력에 항상 같은 출력을 내고, 외부 상태를 바꾸지 않습니다.
* 트랜잭션 abort 후 재처리되어도 안전합니다.
*/
public final class OrderEnricher {
// 부가세율. 실제로는 설정이나 참조 데이터에서 옵니다.
private static final BigDecimal VAT_RATE = new BigDecimal("0.10");
private OrderEnricher() {
}
/**
* 입력 JSON 에 세금 정보를 붙여 반환합니다.
* 파싱 실패는 데이터 자체의 문제이므로 재시도해도 같습니다 →
* 호출자가 격리(DLQ)로 보내야 합니다.
*/
public static String enrich(String rawJson) {
long amount = extractLong(rawJson, "amount");
BigDecimal vat = BigDecimal.valueOf(amount)
.multiply(VAT_RATE)
.setScale(0, RoundingMode.HALF_UP);
BigDecimal total = BigDecimal.valueOf(amount).add(vat);
// 원본을 보존하면서 필드를 덧붙입니다.
// 실제 서비스라면 Jackson 또는 Avro 를 씁니다(예제 7).
return rawJson.substring(0, rawJson.lastIndexOf('}'))
+ ",\"vat\":" + vat.toPlainString()
+ ",\"totalAmount\":" + total.toPlainString()
+ ",\"enrichedBy\":\"eos-pipeline\"}";
}
/** 아주 단순한 숫자 필드 추출기. 형식이 다르면 예외를 던집니다. */
private static long extractLong(String json, String field) {
String needle = "\"" + field + "\":";
int i = json.indexOf(needle);
if (i < 0) {
throw new IllegalArgumentException("필드 없음: " + field + " in " + json);
}
int start = i + needle.length();
int end = start;
while (end < json.length() && (Character.isDigit(json.charAt(end)) || json.charAt(end) == '-')) {
end++;
}
if (start == end) {
throw new IllegalArgumentException("숫자 파싱 실패: " + field + " in " + json);
}
return Long.parseLong(json.substring(start, end));
}
}
트랜잭션 루프 — 이 예제의 본체
package com.example.kafka.eos;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.CommitFailedException;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.clients.consumer.OffsetAndMetadata;
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.common.KafkaException;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.errors.AuthorizationException;
import org.apache.kafka.common.errors.OutOfOrderSequenceException;
import org.apache.kafka.common.errors.ProducerFencedException;
import org.apache.kafka.common.errors.UnsupportedVersionException;
import org.apache.kafka.common.errors.WakeupException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.time.Duration;
import java.util.List;
import java.util.Map;
import java.util.concurrent.atomic.AtomicBoolean;
/**
* consume-transform-produce 파이프라인 (exactly-once).
*
* 한 번의 반복(iteration)이 하나의 트랜잭션입니다.
* beginTransaction()
* → 출력 레코드 send()
* → sendOffsetsToTransaction(입력 오프셋, 컨슈머 그룹 메타데이터)
* commitTransaction()
*
* 커밋이 성공하면 "출력 발행" 과 "입력 오프셋 전진" 이 함께 확정됩니다.
* 어느 단계에서 죽더라도 둘 중 하나만 반영되는 상태가 존재하지 않습니다.
*/
public class EosPipeline implements Runnable, AutoCloseable {
private static final Logger log = LoggerFactory.getLogger(EosPipeline.class);
private static final Duration POLL_TIMEOUT = Duration.ofSeconds(1);
private final Consumer<String, String> consumer;
private final Producer<String, String> producer;
private final String inputTopic;
private final String outputTopic;
private final String dlqTopic;
private final AtomicBoolean running = new AtomicBoolean(true);
private long processed = 0;
private long committedTransactions = 0;
private long abortedTransactions = 0;
private long quarantined = 0;
public EosPipeline(String bootstrapServers,
String groupId,
String transactionalId,
String inputTopic,
String outputTopic,
String dlqTopic) {
this.inputTopic = inputTopic;
this.outputTopic = outputTopic;
this.dlqTopic = dlqTopic;
this.consumer = new KafkaConsumer<>(
EosConfig.consumer(bootstrapServers, groupId, transactionalId + "-consumer"));
this.producer = new KafkaProducer<>(
EosConfig.producer(bootstrapServers, transactionalId));
}
@Override
public void run() {
// initTransactions() 는 딱 한 번 호출합니다.
// 이 호출이 하는 일:
// 1) 트랜잭션 코디네이터를 찾고 PID/에포크를 발급받습니다.
// 2) 같은 transactional.id 로 남아 있던 미완료 트랜잭션을 정리합니다.
// 3) 에포크를 올려 이전 세션(좀비)의 쓰기를 차단합니다.
// max.block.ms 안에 끝나지 않으면 TimeoutException 이 발생합니다.
producer.initTransactions();
log.info("트랜잭션 초기화 완료");
consumer.subscribe(List.of(inputTopic));
try {
while (running.get()) {
ConsumerRecords<String, String> records = consumer.poll(POLL_TIMEOUT);
if (records.isEmpty()) {
continue;
}
processBatchInTransaction(records);
}
} catch (WakeupException e) {
// close() 에서 consumer.wakeup() 을 호출한 정상 종료 경로입니다.
log.info("종료 신호 수신");
} finally {
log.info("파이프라인 종료. {}", stats());
}
}
/**
* poll 로 받은 배치 하나를 하나의 트랜잭션으로 처리합니다.
*/
private void processBatchInTransaction(ConsumerRecords<String, String> records) {
producer.beginTransaction();
try {
for (ConsumerRecord<String, String> record : records) {
String output;
try {
output = OrderEnricher.enrich(record.value());
} catch (RuntimeException dataError) {
// 데이터 자체가 잘못된 경우입니다. 재시도해도 같은 결과이므로
// 트랜잭션을 abort 하지 않고 DLQ 로 보냅니다.
// DLQ 쓰기도 같은 트랜잭션 안이므로 원자성이 유지됩니다.
log.warn("변환 실패 → DLQ. {}-{}@{} key={}",
record.topic(), record.partition(), record.offset(), record.key());
producer.send(new ProducerRecord<>(dlqTopic, record.key(), record.value()));
quarantined++;
continue;
}
// 키를 유지해야 하류의 파티션·순서 보장이 이어집니다.
producer.send(new ProducerRecord<>(outputTopic, record.key(), output));
processed++;
}
// 커밋할 오프셋은 "다음에 읽을 위치" 입니다. 즉 마지막 오프셋 + 1.
// 4.x 는 ConsumerRecords#nextOffsets() 로 리더 에포크까지 포함한
// 정확한 값을 제공합니다. 직접 offset+1 을 계산하면
// 리더 에포크 정보가 빠져 fencing 이 약해집니다.
Map<TopicPartition, OffsetAndMetadata> offsets = records.nextOffsets();
// groupMetadata() 를 넘기면 브로커가 컨슈머 그룹 세대(generation)까지
// 검증해 더 강한 fencing 을 제공합니다(브로커 2.5+ 필요).
// ConsumerGroupMetadata 를 직접 만들어 넘기면 이 이점을 잃습니다.
producer.sendOffsetsToTransaction(offsets, consumer.groupMetadata());
producer.commitTransaction();
committedTransactions++;
} catch (ProducerFencedException
| OutOfOrderSequenceException
| AuthorizationException
| UnsupportedVersionException fatal) {
// 복구 불가능한 오류입니다.
// - ProducerFencedException: 같은 transactional.id 로 더 새로운 인스턴스가 떴습니다.
// 이 프로세스는 좀비이므로 즉시 죽는 것이 정답입니다.
// - OutOfOrderSequenceException: 시퀀스가 깨졌습니다. 상태를 신뢰할 수 없습니다.
// - AuthorizationException: transactional.id 또는 group.id 권한이 없습니다.
// abortTransaction() 을 호출해도 실패하므로 시도하지 않습니다.
log.error("복구 불가 트랜잭션 오류 — 프로세스를 종료합니다", fatal);
running.set(false);
throw fatal;
} catch (CommitFailedException e) {
// 리밸런스로 파티션을 잃어 오프셋 커밋이 불가능해진 경우입니다.
// 공식 문서가 "abort 로 처리하라" 고 명시합니다.
log.warn("오프셋 커밋 실패(리밸런스 추정) — 트랜잭션을 abort 합니다", e);
abortTransactionQuietly();
} catch (KafkaException e) {
// 그 밖의 오류는 abort 후 다음 반복에서 같은 입력을 다시 읽습니다.
// abort 하면 이 트랜잭션이 쓴 출력 레코드는 read_committed 컨슈머에게
// 보이지 않으므로 중복이 생기지 않습니다.
log.error("트랜잭션 실패 — abort 후 재시도합니다", e);
abortTransactionQuietly();
}
}
private void abortTransactionQuietly() {
try {
producer.abortTransaction();
abortedTransactions++;
} catch (KafkaException e) {
// abort 조차 실패하면 프로듀서 상태를 신뢰할 수 없습니다.
log.error("abortTransaction 실패 — 프로세스를 종료합니다", e);
running.set(false);
}
}
public String stats() {
return "processed=%d, committed=%d, aborted=%d, quarantined=%d"
.formatted(processed, committedTransactions, abortedTransactions, quarantined);
}
/** 다른 스레드(shutdown hook)에서 호출합니다. */
@Override
public void close() {
running.set(false);
// poll() 에서 블록 중인 컨슈머를 깨웁니다. WakeupException 이 던져집니다.
consumer.wakeup();
}
/** run() 스레드가 끝난 뒤에 호출합니다. */
public void closeClients() {
try {
// 진행 중인 트랜잭션이 있으면 abort 됩니다(close 가 정리합니다).
producer.close(Duration.ofSeconds(30));
} finally {
// 4.3에서 consumer.close(Duration) 은 deprecated 입니다.
// 타임아웃을 주려면 CloseOptions 를 씁니다.
consumer.close(org.apache.kafka.clients.consumer.CloseOptions
.timeout(Duration.ofSeconds(30)));
}
}
}
진입점 — transactional.id를 어떻게 정하는가
package com.example.kafka.eos;
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 GROUP_ID = "orders-enricher";
private static final String INPUT = "orders.raw";
private static final String OUTPUT = "orders.enriched";
private static final String DLQ = "orders.enriched.DLT";
public static void main(String[] args) throws InterruptedException {
// 인스턴스 ID: 환경변수 → 인자 → 기본값 순으로 결정합니다.
// 절대 랜덤으로 만들지 않습니다(재시작 시 fencing 이 깨집니다).
String instanceId = System.getenv("INSTANCE_ID");
if (instanceId == null || instanceId.isBlank()) {
instanceId = args.length > 0 ? args[0] : "0";
}
String transactionalId = "orders-enricher-tx-" + instanceId;
log.info("transactional.id = {}", transactionalId);
EosPipeline pipeline = new EosPipeline(
BOOTSTRAP, GROUP_ID, transactionalId, INPUT, OUTPUT, DLQ);
Thread worker = new Thread(pipeline, "eos-pipeline-" + instanceId);
Runtime.getRuntime().addShutdownHook(new Thread(() -> {
log.info("SIGTERM 수신 — 정상 종료를 시작합니다");
pipeline.close(); // poll() 을 깨웁니다
try {
worker.join(60_000); // 진행 중인 트랜잭션이 끝날 시간을 줍니다
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
pipeline.closeClients();
}, "shutdown-hook"));
worker.start();
worker.join();
}
}
실행 방법
# 0. 예제 1의 클러스터가 떠 있어야 합니다.
cd kafka-lab && docker compose ps
# 1. DLQ 토픽을 추가로 만듭니다(예제 1의 create-topics.sh 에는 없습니다).
./kcli kafka-topics.sh --create --if-not-exists \
--topic orders.enriched.DLT --partitions 3 --replication-factor 3 \
--config min.insync.replicas=2
# 2. 파이프라인 기동 (인스턴스 0)
cd ../eos-pipeline
mvn -q clean package
INSTANCE_ID=0 mvn -q exec:java
# 3. 다른 터미널에서 입력 이벤트를 넣습니다.
cd kafka-lab
for i in $(seq 1 20); do
echo "ORD-$i:{\"orderId\":\"ORD-$i\",\"amount\":$((10000 * i)),\"currency\":\"KRW\"}"
done | ./kcli kafka-console-producer.sh --topic orders.raw \
--property parse.key=true --property key.separator=:
# 4. 잘못된 데이터 1건을 넣어 DLQ 경로도 확인합니다.
echo 'ORD-BAD:{"orderId":"ORD-BAD","currency":"KRW"}' \
| ./kcli kafka-console-producer.sh --topic orders.raw \
--property parse.key=true --property key.separator=:
[main] INFO com.example.kafka.eos.Main - transactional.id = orders-enricher-tx-0
[eos-pipeline-0] INFO o.a.k.c.p.internals.TransactionManager - [Producer clientId=orders-enricher-tx-0, transactionalId=orders-enricher-tx-0] ProducerId set to 5000 with epoch 0
[eos-pipeline-0] INFO com.example.kafka.eos.EosPipeline - 트랜잭션 초기화 완료
[eos-pipeline-0] WARN com.example.kafka.eos.EosPipeline - 변환 실패 → DLQ. orders.raw-4@5 key=ORD-BAD
검증 방법
1. 출력이 정확히 한 번씩 있는가
read_committed로 읽어 키별 중복 확인./kcli kafka-console-consumer.sh --topic orders.enriched --from-beginning \
--timeout-ms 10000 \
--consumer-property isolation.level=read_committed \
--property print.key=true 2>/dev/null | sort | uniq -c | sort -rn | head
1 ORD-1 {"orderId":"ORD-1","amount":10000,"currency":"KRW","vat":1000,"totalAmount":11000,"enrichedBy":"eos-pipeline"}
1 ORD-10 {"orderId":"ORD-10","amount":100000,"currency":"KRW","vat":10000,"totalAmount":110000,"enrichedBy":"eos-pipeline"}
1 ORD-11 {"orderId":"ORD-11","amount":110000,"currency":"KRW","vat":11000,"totalAmount":121000,"enrichedBy":"eos-pipeline"}
2. read_committed와 read_uncommitted의 차이를 눈으로 본다
이것이 이 예제의 가장 교육적인 검증입니다. 트랜잭션 마커(commit/abort 표시)는 로그에 실제 레코드로 기록되므로, 두 격리 수준의 건수와 오프셋이 다르게 보입니다.
echo -n 'read_committed : '
./kcli kafka-console-consumer.sh --topic orders.enriched --from-beginning \
--timeout-ms 8000 --consumer-property isolation.level=read_committed \
2>/dev/null | wc -l
echo -n 'read_uncommitted: '
./kcli kafka-console-consumer.sh --topic orders.enriched --from-beginning \
--timeout-ms 8000 --consumer-property isolation.level=read_uncommitted \
2>/dev/null | wc -l
# 파티션별 log-end-offset. 트랜잭션 마커가 오프셋을 소비하므로
# 이 합계는 실제 레코드 수보다 큽니다.
./kcli kafka-get-offsets.sh --topic orders.enriched --time -1
정상 종료된 파이프라인에서는 두 격리 수준의 레코드 수가 같습니다
(abort된 트랜잭션이 없으므로). 하지만
kafka-get-offsets.sh의 합계는 레코드 수보다 큽니다.
차이가 바로 트랜잭션 커밋 마커가 차지한 오프셋입니다.
"오프셋 합계와 메시지 수가 안 맞는다"는 신고의 흔한 원인이 이것입니다.
read_committed와 LSO — 진행 중인 트랜잭션 때문에
그 뒤의 커밋된 레코드까지 보이지 않게 되는 상황
3. abort된 트랜잭션이 하류에 보이지 않는가
파이프라인을 처리 도중에 강제 종료해 미완료 트랜잭션을 만들고,
read_committed에는 나타나지 않는지 확인합니다.
# 입력을 대량으로 넣어 처리 중 상태를 만듭니다.
for i in $(seq 100 400); do
echo "ORD-$i:{\"orderId\":\"ORD-$i\",\"amount\":$((100 * i)),\"currency\":\"KRW\"}"
done | ./kcli kafka-console-producer.sh --topic orders.raw \
--property parse.key=true --property key.separator=:
# 파이프라인 프로세스를 SIGKILL 로 죽입니다(정상 종료 경로를 건너뜁니다).
pkill -9 -f 'com.example.kafka.eos.Main'
# 같은 INSTANCE_ID 로 재시작 — 같은 transactional.id 를 쓰므로
# 브로커가 이전 세션의 미완료 트랜잭션을 정리하고 에포크를 올립니다.
cd ../eos-pipeline && INSTANCE_ID=0 mvn -q exec:java
[eos-pipeline-0] INFO o.a.k.c.p.internals.TransactionManager - [Producer clientId=orders-enricher-tx-0, transactionalId=orders-enricher-tx-0] ProducerId set to 5000 with epoch 1
epoch가 0에서 1로 올라간 것이 fencing이 동작한 증거입니다.
이제 옛 프로세스(있었다면)가 같은 transactional.id로 쓰려 하면
ProducerFencedException을 받고 죽습니다.
그리고 위 검증 1을 다시 실행하면
여전히 모든 키의 카운트가 1입니다 — 크래시했는데도 중복이 없습니다.
4. 오프셋과 출력이 같이 전진했는가
./kcli kafka-consumer-groups.sh --describe --group orders-enricher
GROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG CONSUMER-ID HOST CLIENT-ID
orders-enricher orders.raw 0 54 54 0 orders-...-0-... /172.18.0.1 orders-enricher-tx-0-consumer
orders-enricher orders.raw 1 51 51 0 orders-...-0-... /172.18.0.1 orders-enricher-tx-0-consumer
LAG이 0이면 입력을 모두 처리했다는 뜻입니다.
여기서 CURRENT-OFFSET은 트랜잭션과 함께 커밋된 값입니다.
즉 이 오프셋이 전진했다는 것은 그 구간의 출력이 반드시 커밋되었다는 것과 동의어입니다.
EOS의 경계 — 가장 흔한 오해
Kafka 트랜잭션이 롤백할 수 있는 것은 Kafka에 쓴 것뿐입니다. 트랜잭션 안에서 DB에 INSERT하거나 결제 API를 호출하면 그 부수효과는 abort되어도 남습니다.
트랜잭션 안에서 외부 시스템을 호출합니다.
abortTransaction()이 실행되면 Kafka 출력은 사라지지만
DB 행과 이메일은 남습니다. 재처리 시 두 번째 이메일이 발송됩니다.
producer.beginTransaction();
for (var record : records) {
var order = parse(record.value());
orderRepository.insert(order); // ← 롤백되지 않습니다
mailService.sendConfirmation(order);// ← 취소되지 않습니다
producer.send(new ProducerRecord<>(
outputTopic, record.key(), toJson(order)));
}
producer.sendOffsetsToTransaction(
records.nextOffsets(), consumer.groupMetadata());
producer.commitTransaction();
트랜잭션 안에서는 Kafka 쓰기만 합니다. 외부 시스템 반영은 별도 컨슈머가 자체 멱등 키로 처리합니다. 같은 이벤트를 두 번 받아도 결과가 같습니다.
producer.beginTransaction();
for (var record : records) {
// 순수 변환만 수행합니다.
producer.send(new ProducerRecord<>(
outputTopic, record.key(),
OrderEnricher.enrich(record.value())));
}
producer.sendOffsetsToTransaction(
records.nextOffsets(), consumer.groupMetadata());
producer.commitTransaction();
// UNIQUE(event_id) 제약으로 중복 삽입을 DB 가 막습니다.
// 같은 이벤트를 두 번 받아도 결과가 같습니다 = 멱등.
for (var record : records) {
var event = parse(record.value());
int inserted = jdbc.update("""
INSERT INTO settlement (event_id, order_id, total_amount)
VALUES (?, ?, ?)
ON CONFLICT (event_id) DO NOTHING
""", event.id(), event.orderId(), event.totalAmount());
if (inserted == 0) {
log.debug("이미 처리된 이벤트 — 건너뜁니다 id={}", event.id());
}
}
consumer.commitSync();
이 예제에서 쓴 설정
| 설정 | 기본값 | 이 예제 값 | 이유 |
|---|---|---|---|
transactional.idproducer |
null |
orders-enricher-tx-{instanceId} |
설정하면 트랜잭션 API가 활성화되고 멱등성이 강제됩니다. 재시작 후에도 같아야 fencing이 동작합니다 |
transaction.timeout.msproducer |
60000 |
60000 |
한 트랜잭션의 최대 지속 시간. 브로커의 transaction.max.timeout.ms(기본 900000)를 넘길 수 없습니다 |
enable.auto.commitconsumer |
true |
false |
필수. 자동 커밋이 켜져 있으면 트랜잭션 밖에서 오프셋이 전진해 원자성이 깨집니다 |
isolation.levelconsumer |
read_uncommitted |
read_committed |
기본값이 read_uncommitted라는 점이 함정입니다. 하류 컨슈머도 반드시 바꿔야 EOS가 완성됩니다 |
max.poll.recordsconsumer |
500 |
200 |
배치 처리 시간이 transaction.timeout.ms를 넘으면 커밋이 실패합니다 |
transaction.state.log.replication.factorbroker |
3 |
3 |
__transaction_state의 복제 계수. 1이면 트랜잭션 상태가 단일 사본이 되어 보장이 무의미합니다 |
transaction.state.log.min.isrbroker |
2 |
2 |
트랜잭션 상태 로그의 내구성 하한 |
transactional.id.expiration.msbroker |
604800000 |
기본값 유지 | 7일 동안 쓰이지 않은 transactional.id는 코디네이터가 만료시킵니다 |
processing.guaranteestreams |
at_least_once |
(이 예제는 Streams를 쓰지 않음) | Streams에서 같은 보장을 원하면 exactly_once_v2. 유효값은 이 둘뿐입니다 |
프로덕션 고려사항
| 항목 | 이 예제 | 프로덕션 |
|---|---|---|
transactional.id 부여 |
환경변수 또는 인자 | K8s StatefulSet의 안정적 파드 이름에서 파생. Deployment(랜덤 파드 이름)로 운영하면 매 배포마다 새 ID가 생겨 fencing이 무력화되고 코디네이터에 좀비 ID가 누적됩니다 |
| 처리량 | 배치당 1트랜잭션 | 트랜잭션 커밋에는 왕복 비용이 있습니다. 배치를 너무 작게 잡으면 처리량이 급락합니다. max.poll.records와 transaction.timeout.ms 사이에서 균형점을 찾아야 합니다 |
| 하류 컨슈머 | CLI로 확인 | 모든 하류 컨슈머가 read_committed여야 합니다. 하나라도 기본값(read_uncommitted)이면 abort된 레코드를 읽어 EOS가 무너집니다. 팀 간 계약으로 관리해야 합니다 |
| 지연 | 측정하지 않음 | read_committed 컨슈머는 LSO까지만 읽습니다. 긴 트랜잭션이 하나 열려 있으면 그 뒤의 커밋된 레코드까지 보이지 않아 end-to-end 지연이 트랜잭션 길이만큼 늘어납니다 |
| 종료 처리 | shutdown hook + join(60s) |
terminationGracePeriodSeconds를 transaction.timeout.ms보다 크게 잡습니다. 짧으면 SIGKILL로 미완료 트랜잭션이 남아 하류 컨슈머가 LSO에서 멈춥니다 |
| 모니터링 | 애플리케이션 카운터 | abort 비율, 트랜잭션 커밋 지연, 하류 컨슈머 lag을 함께 봅니다. lag이 계단식으로 튀면 긴 트랜잭션이 LSO를 막고 있다는 신호입니다(예제 10) |
| Streams 대안 | 클라이언트 API 직접 사용 | 단순한 consume-transform-produce라면 Kafka Streams + processing.guarantee=exactly_once_v2가 코드가 훨씬 짧고 실수 여지가 적습니다(예제 9). 세밀한 제어가 필요할 때만 직접 씁니다 |
자주 하는 실수
관련 케이스 스터디
이어서 볼 곳
공식 문서 출처
- KafkaProducer Javadoc — 트랜잭션 API 사용 예시와 예외 처리 패턴(
ProducerFencedException/OutOfOrderSequenceException/AuthorizationException은 복구 불가, 그 외KafkaException은 abort),transactional.id의 목적, 트랜잭션 토픽에 RF≥3·min.insync.replicas=2권고 sendOffsetsToTransactionJavadoc — 커밋 오프셋은 "다음에 읽을 위치"이며ConsumerRecords#nextOffsets()를 쓰라는 안내,groupMetadata()가 더 강한 fencing을 제공한다는 서술,enable.auto.commit=false이고 수동 커밋도 하지 말라는 요구,CommitFailedException은 abort로 처리하라는 안내- Producer Configs —
transactional.id=null,transaction.timeout.ms=60000 - Consumer Configs —
isolation.level=read_uncommitted,enable.auto.commit=true,auto.commit.interval.ms=5000,max.poll.records=500 - Broker Configs —
transaction.state.log.replication.factor=3,transaction.state.log.min.isr=2,transaction.max.timeout.ms=900000,transactional.id.expiration.ms=604800000 - Kafka Streams Configs —
processing.guarantee기본값at_least_once, 유효값at_least_once/exactly_once_v2, EOS일 때commit.interval.ms가100이 된다는 서술 - Transaction Protocol — 코디네이터·
__transaction_state·마커의 동작