학습 목표

시나리오

주문 수집 서비스가 orders.raw 토픽에 원본 주문을 넣습니다. 중간 파이프라인이 이 원본을 읽어 통화 정규화와 세금 계산을 하고 orders.enriched 토픽에 씁니다. 하류의 정산 시스템은 orders.enriched만 봅니다.

문제는 이 파이프라인이 배포·크래시·리밸런스로 재시작될 때입니다. at-least-once로 짜면 같은 주문의 보강 결과가 두 번 발행되어 정산 금액이 두 배가 됩니다. at-most-once로 짜면 누락됩니다. 정산 시스템에 멱등 처리를 넣는 것이 정석이지만, 레거시라 손댈 수 없다는 제약이 있습니다.

입력과 출력이 모두 Kafka 토픽이므로 이 구간은 EOS의 적용 대상입니다. Kafka 트랜잭션으로 처리합니다.

아키텍처

consume-transform-produce — 오프셋 커밋이 트랜잭션 안에 들어갑니다 입력 토픽에서 읽어 변환한 뒤 출력 토픽으로 쓰는 패턴을 그린 그림입니다. KafkaConsumer 는 isolation.level=read_committed, enable.auto.commit=false 로 설정합니다. 변환 결과를 트랜잭션 프로듀서로 출력 토픽에 쓰고, 같은 트랜잭션 안에서 sendOffsetsToTransaction 으로 컨슈머의 오프셋을 __consumer_offsets 에 씁니다. 점선으로 표시한 트랜잭션 경계 안에 출력 레코드와 오프셋 커밋이 함께 들어 있는 것이 핵심입니다. 그래서 커밋하면 둘 다 반영되고, abort 하면 둘 다 무효가 되어 같은 레코드를 다시 처리합니다. 컨슈머의 commitSync 나 commitAsync 를 따로 호출하면 원자성이 깨지므로 호출하지 않습니다. consume-transform-produce — 출력과 오프셋을 한 트랜잭션으로 묶습니다 하나의 트랜잭션 (원자적) 입력 토픽 orders KafkaConsumer read_committed auto.commit=false 변환 처리 transform 트랜잭션 프로듀서 transactional.id 컨슈머 1:1 로 두기 출력 토픽 orders-enriched sendOffsets ToTransaction() __consumer _offsets 커밋하면 출력 레코드와 오프셋이 동시에 확정됩니다. abort 하면 둘 다 무효 → 같은 구간을 재처리. 하지 말 것 consumer.commitSync() 따로 커밋하면 원자성이 깨집니다. 오프셋은 프로듀서가 컨슈머 그룹을 대신해 씁니다 — 그래서 consumer.groupMetadata() 를 함께 넘깁니다. abort 후에는 오프셋을 자동으로 되감지 않습니다. 애플리케이션이 커밋된 오프셋을 다시 읽어 위치를 되돌려야 합니다.
consume-transform-produce — 출력 레코드와 입력 오프셋 커밋이 하나의 트랜잭션에 함께 들어가는 구조
트랜잭션 흐름 — initTransactions 부터 commit 또는 abort 까지 트랜잭션 프로듀서의 API 호출 순서와 각 호출이 트랜잭션 코디네이터에서 무엇을 하는지, 그 결과가 어디에 기록되는지를 6단계로 보여줍니다. 1단계 initTransactions 는 코디네이터를 찾아 PID 와 epoch 를 발급받고 이전 인스턴스의 미완 트랜잭션을 정리합니다. 2단계 beginTransaction 은 클라이언트 로컬 상태 전환이라 브로커와 통신하지 않습니다. 3단계 첫 send 전에 AddPartitionsToTxn 으로 대상 파티션이 __transaction_state 에 등록되고 레코드가 데이터 파티션에 기록됩니다. 4단계 sendOffsetsToTransaction 은 AddOffsetsToTxn 과 TxnOffsetCommit 으로 컨슈머 오프셋을 같은 트랜잭션 안에서 __consumer_offsets 에 씁니다. 5단계 commitTransaction 은 EndTxn 요청으로 코디네이터가 PREPARE_COMMIT 을 기록하고 관련된 모든 파티션에 커밋 마커를 쓴 뒤 COMPLETE_COMMIT 으로 끝냅니다. 6단계 abortTransaction 은 같은 경로로 abort 마커를 쓰며, read_committed 컨슈머가 그 레코드를 건너뜁니다. 코디네이터의 상태는 내부 토픽 __transaction_state 에 저장됩니다. 트랜잭션 흐름 — 프로듀서 호출 · 코디네이터 동작 · 기록 위치 프로듀서 API 호출 트랜잭션 코디네이터 기록 위치 1 initTransactions() 앱 시작 시 한 번만 코디네이터 찾기 → PID · epoch 발급, 미완 정리 __transaction_state Empty 상태로 초기화 2 beginTransaction() 트랜잭션마다 호출 브로커와 통신하지 않습니다 — 클라이언트 로컬 상태만 바뀝니다 코디네이터에는 첫 send() 시점에 처음 기록됩니다 3 send(record) 여러 토픽·파티션 가능 AddPartitionsToTxn 참여 파티션 목록을 기록 대상 토픽 파티션 이미 로그에 있지만 미확정 4 sendOffsets ToTransaction() AddOffsetsToTxn TxnOffsetCommit __consumer_offsets 같은 트랜잭션 안에서 5 commitTransaction() 여기서 원자성이 확정 EndTxn PREPARE → 마커 → 완료 모든 참여 파티션에 COMMIT 마커 기록 6 abortTransaction() 예외 처리 경로 EndTxn (ABORT) 오프셋 커밋도 함께 무효 ABORT 마커 기록 read_committed 가 건너뜀 레코드는 기록된 뒤에 마커로 확정됩니다. 그래서 abort 된 레코드도 로그에는 남아 있고, 컨슈머가 걸러냅니다. transaction.timeout.ms 기본 60000 · transactional.id.expiration.ms 기본 604800000 transaction.state.log.replication.factor 기본 3 · transaction.state.log.min.isr 기본 2
트랜잭션 흐름 — initTransactions부터 commitTransaction까지, 트랜잭션 코디네이터와 __transaction_state의 역할

핵심 메커니즘 세 가지

EOS를 구성하는 세 부품
부품 담당하는 문제 없으면 무슨 일이 생기는가
transactional.id + 에포크 재시작·좀비 프로세스 차단(fencing) 네트워크 분리로 살아남은 옛 인스턴스가 계속 써서 중복이 생깁니다
sendOffsetsToTransaction 출력 쓰기와 입력 오프셋 커밋의 원자성 둘 사이에서 죽으면 중복 또는 유실이 됩니다
isolation.level=read_committed 하류 컨슈머가 미완료·중단된 트랜잭션을 보지 않게 함 abort된 레코드까지 읽어 정산이 틀어집니다

사전 요구사항

전체 코드

디렉터리 구조
eos-pipeline/
├── pom.xml
└── src/main/java/com/example/kafka/eos/
    ├── EosConfig.java          # 프로듀서/컨슈머 설정
    ├── OrderEnricher.java      # 순수 변환 로직 (부수효과 없음)
    ├── EosPipeline.java        # 트랜잭션 루프 본체
    └── Main.java               # 진입점 (인스턴스 ID 로 transactional.id 생성)

pom.xml

eos-pipeline/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>

설정 — 트랜잭션에 필요한 값만 정확히

src/main/java/com/example/kafka/eos/EosConfig.java
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의 경계에서 다룹니다.

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

트랜잭션 루프 — 이 예제의 본체

src/main/java/com/example/kafka/eos/EosPipeline.java
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를 어떻게 정하는가

src/main/java/com/example/kafka/eos/Main.java
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이어야 합니다
      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_committedread_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 — 진행 중 트랜잭션이 뒤쪽 커밋 메시지까지 가립니다 파티션 로그의 오프셋 0부터 11까지를 12칸으로 그린 그림입니다. 오프셋 0과 1은 일반 메시지, 2와 3은 트랜잭션 T1 의 레코드, 4는 T1 의 커밋 마커입니다. 오프셋 5는 아직 진행 중인 트랜잭션 T2 의 첫 레코드이고, 6과 7은 트랜잭션 T3 의 레코드, 8은 T3 의 커밋 마커, 9와 11은 트랜잭션과 무관한 일반 메시지, 10은 T2 의 두 번째 레코드입니다. LSO 즉 last stable offset 은 진행 중인 가장 앞선 트랜잭션의 첫 오프셋인 5 입니다. isolation.level 이 기본값 read_uncommitted 인 컨슈머는 high watermark 까지 전부 읽습니다. read_committed 컨슈머는 LSO 보다 작은 오프셋만 반환하므로 0부터 4까지만 읽습니다. 그 결과 T3 가 이미 커밋을 끝낸 6, 7, 8 과 트랜잭션과 아무 상관 없는 9, 11 까지 보이지 않습니다. 이것이 read_committed 가 지연을 만드는 이유이며, 진행 중 트랜잭션이 길수록 지연이 커집니다. read_committed 와 LSO — 진행 중 트랜잭션 하나가 그 뒤 전부를 막습니다 리더 로그 orders-0 0 일반 보임 1 일반 보임 2 T1 커밋 3 T1 커밋 4 마커 T1 끝 5 T2 진행중 6 T3 커밋 7 T3 커밋 8 마커 T3 끝 9 일반 막힘 10 T2 진행중 11 일반 막힘 HW 12 LSO = 5 진행 중 T2 의 첫 오프셋 read_uncommitted (기본값) 0 ~ 11 전부 반환 — abort 된 레코드까지 보입니다 read_committed 0 ~ 4 만 반환 LSO 뒤 — 7건 전부 보류 오프셋 6·7·8 은 T3 가 이미 커밋을 끝낸 레코드이고, 9·11 은 트랜잭션과 무관한 일반 메시지입니다. 그런데도 오프셋 5의 T2 가 끝나지 않아 전부 보류됩니다. 긴 트랜잭션 = 컨슈머 지연이 되는 이유입니다. read_committed 로도 보임 진행 중 트랜잭션 T2 의 레코드 LSO 뒤라서 보이지 않음 isolation.level 기본값은 read_uncommitted 입니다. read_committed 에서 seekToEnd() 는 HW 가 아니라 LSO 를 돌려줍니다.
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의 경계 — 가장 흔한 오해

EOS 의 경계 — Kafka 내부는 exactly-once, 외부 시스템은 아닙니다 exactly-once semantics 가 성립하는 범위와 성립하지 않는 범위를 굵은 경계선으로 나눈 그림입니다. 경계선 위쪽은 Kafka 내부입니다. 입력 토픽에서 읽어 처리한 뒤 출력 토픽에 쓰고 컨슈머 오프셋을 __consumer_offsets 에 쓰는 것까지가 하나의 트랜잭션이므로 원자적으로 확정됩니다. 경계선 아래쪽은 외부 시스템입니다. 애플리케이션이 처리 중에 호출하는 외부 데이터베이스 INSERT 나 UPDATE, 결제 API 의 HTTP POST, 이메일이나 SMS 발송은 Kafka 트랜잭션에 참여할 수 없습니다. 경계를 통과하는 화살표에 X 표시를 한 것은 그 지점에서 보장이 끊긴다는 뜻입니다. 따라서 트랜잭션이 abort 되거나 재처리가 일어나면 Kafka 쪽 결과는 되돌아가지만 이미 실행된 외부 작업은 되돌아가지 않아 중복이 남습니다. EOS 를 켜면 데이터베이스 중복까지 막힌다는 것은 오해입니다. 해결책은 컨슈머 또는 싱크 쪽의 멱등성입니다. 업서트, 유니크 제약, 처리 기록 테이블로 이미 처리한 레코드를 걸러내야 합니다. EOS 의 경계 — 어디까지 Kafka 가 보장하는가 EOS 성립 — Kafka 토픽 → Kafka 토픽 (트랜잭션 원자성) 입력 토픽 orders 애플리케이션 (Streams 또는 직접 구현) read_committed · auto.commit=false transactional.id 설정 Streams: exactly_once_v2 출력 토픽 — 커밋 마커로 확정 __consumer_offsets 같은 트랜잭션 안에서 커밋 EOS 경계 — Kafka 의 보장은 여기서 끊깁니다 처리 중 외부 호출 — 트랜잭션에 참여하지 못합니다 EOS 성립 안 함 — 외부 시스템 (커밋도 롤백도 Kafka 와 무관) 외부 DB — INSERT / UPDATE 결제 API — HTTP POST 이메일 · SMS 발송 “EOS 를 켰으니 DB 중복도 막힌다”는 오해입니다. abort 나 재처리가 일어나면 Kafka 쪽은 되돌아가지만, 이미 실행된 외부 작업은 그대로 남습니다. 해결: 경계 밖에서는 컨슈머(싱크) 스스로 멱등성을 확보해야 합니다 업서트(UPSERT) · 유니크 제약 · 처리 기록 테이블(메시지 키 + 오프셋)로 재처리를 걸러냅니다. 공식 문서: 외부 시스템에 쓸 때의 한계는 컨슈머의 위치와 실제 저장된 출력을 함께 맞춰야 한다는 점입니다. 다른 목적지 시스템의 exactly-once 는 그 시스템의 협조가 필요합니다.
EOS의 경계 — Kafka 내부는 원자적이지만 외부 DB·API 호출은 트랜잭션에 포함되지 않습니다

Kafka 트랜잭션이 롤백할 수 있는 것은 Kafka에 쓴 것뿐입니다. 트랜잭션 안에서 DB에 INSERT하거나 결제 API를 호출하면 그 부수효과는 abort되어도 남습니다.

트랜잭션 안에서 외부 시스템을 호출합니다. abortTransaction()이 실행되면 Kafka 출력은 사라지지만 DB 행과 이메일은 남습니다. 재처리 시 두 번째 이메일이 발송됩니다.

EosPipeline.java
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 쓰기만 합니다. 외부 시스템 반영은 별도 컨슈머가 자체 멱등 키로 처리합니다. 같은 이벤트를 두 번 받아도 결과가 같습니다.

EosPipeline.java (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();
SettlementConsumer.java (외부 시스템 반영)
// 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();

이 예제에서 쓴 설정

Apache Kafka 4.3 기준
설정 기본값 이 예제 값 이유
transactional.id
producer
null orders-enricher-tx-{instanceId} 설정하면 트랜잭션 API가 활성화되고 멱등성이 강제됩니다. 재시작 후에도 같아야 fencing이 동작합니다
transaction.timeout.ms
producer
60000 60000 한 트랜잭션의 최대 지속 시간. 브로커의 transaction.max.timeout.ms(기본 900000)를 넘길 수 없습니다
enable.auto.commit
consumer
true false 필수. 자동 커밋이 켜져 있으면 트랜잭션 밖에서 오프셋이 전진해 원자성이 깨집니다
isolation.level
consumer
read_uncommitted read_committed 기본값이 read_uncommitted라는 점이 함정입니다. 하류 컨슈머도 반드시 바꿔야 EOS가 완성됩니다
max.poll.records
consumer
500 200 배치 처리 시간이 transaction.timeout.ms를 넘으면 커밋이 실패합니다
transaction.state.log.replication.factor
broker
3 3 __transaction_state의 복제 계수. 1이면 트랜잭션 상태가 단일 사본이 되어 보장이 무의미합니다
transaction.state.log.min.isr
broker
2 2 트랜잭션 상태 로그의 내구성 하한
transactional.id.expiration.ms
broker
604800000 기본값 유지 7일 동안 쓰이지 않은 transactional.id는 코디네이터가 만료시킵니다
processing.guarantee
streams
at_least_once (이 예제는 Streams를 쓰지 않음) Streams에서 같은 보장을 원하면 exactly_once_v2. 유효값은 이 둘뿐입니다

프로덕션 고려사항

로컬 예제와 프로덕션의 차이
항목이 예제프로덕션
transactional.id 부여 환경변수 또는 인자 K8s StatefulSet의 안정적 파드 이름에서 파생. Deployment(랜덤 파드 이름)로 운영하면 매 배포마다 새 ID가 생겨 fencing이 무력화되고 코디네이터에 좀비 ID가 누적됩니다
처리량 배치당 1트랜잭션 트랜잭션 커밋에는 왕복 비용이 있습니다. 배치를 너무 작게 잡으면 처리량이 급락합니다. max.poll.recordstransaction.timeout.ms 사이에서 균형점을 찾아야 합니다
하류 컨슈머 CLI로 확인 모든 하류 컨슈머가 read_committed여야 합니다. 하나라도 기본값(read_uncommitted)이면 abort된 레코드를 읽어 EOS가 무너집니다. 팀 간 계약으로 관리해야 합니다
지연 측정하지 않음 read_committed 컨슈머는 LSO까지만 읽습니다. 긴 트랜잭션이 하나 열려 있으면 그 뒤의 커밋된 레코드까지 보이지 않아 end-to-end 지연이 트랜잭션 길이만큼 늘어납니다
종료 처리 shutdown hook + join(60s) terminationGracePeriodSecondstransaction.timeout.ms보다 크게 잡습니다. 짧으면 SIGKILL로 미완료 트랜잭션이 남아 하류 컨슈머가 LSO에서 멈춥니다
모니터링 애플리케이션 카운터 abort 비율, 트랜잭션 커밋 지연, 하류 컨슈머 lag을 함께 봅니다. lag이 계단식으로 튀면 긴 트랜잭션이 LSO를 막고 있다는 신호입니다(예제 10)
Streams 대안 클라이언트 API 직접 사용 단순한 consume-transform-produce라면 Kafka Streams + processing.guarantee=exactly_once_v2가 코드가 훨씬 짧고 실수 여지가 적습니다(예제 9). 세밀한 제어가 필요할 때만 직접 씁니다

자주 하는 실수

이어서 볼 곳

공식 문서 출처