Step 13 — 실전 패턴과 최종 프로젝트

학습 목표

  • at-least-once 위에서 멱등성으로 "정확히 한 번"을 만드는 세 축(멱등 컨슈머 / 원자적 발행 / 순서 보장)을 구분한다
  • processed_message 처리 이력 테이블로 멱등 컨슈머를 구현하고, 메시지 ID 후보 3가지를 비교해 하나를 고른다
  • SELECTINSERT 하는 멱등 체크가 왜 경합에서 깨지는지 두 스레드로 재현한다
  • OrderService + OutboxRelayTransactional Outbox 를 구현하고, FOR UPDATE SKIP LOCKED 로 다중 인스턴스 안전성을 확인한다
  • 순서 보장이 깨지는 5가지 경로를 표로 정리하고, @Async 리스너가 순서와 커밋을 동시에 깨뜨리는 것을 로그로 확인한다
  • 주문 서비스 전체를 조립해 검증 시나리오 4개를 실행하고, 배포 전 체크리스트 20항목으로 Step 01~12 를 회수한다

선행 스텝: Step 12 — 관측성과 운영 예상 소요: 150분


13-0. 실습 준비

이 스텝은 Kafka 와 MySQL 을 동시에 씁니다. docker compose ps 로 둘 다 (healthy) 인지 먼저 확인하세요. MySQL 별칭을 하나 더 만들어 둡니다. 이 스텝 내내 씁니다.

alias mq='docker exec -i learn-kafka-mysql mysql -ulearner -plearn1234 orderdb -t -e'

mq "SHOW TABLES;"        # inventory / orders / outbox_event / processed_message 4개
kt --create --topic orders.outbox     --partitions 3 --replication-factor 1
kt --create --topic orders.outbox.DLT --partitions 3 --replication-factor 1

이 스텝은 REST 엔드포인트를 씁니다. build.gradle 에 웹 스타터를 추가하세요. (Step 12 에서 Actuator 웹 엔드포인트를 열었다면 이미 있습니다.)

implementation 'org.springframework.boot:spring-boot-starter-web'

💡 실습 도중 상태가 꼬이면 DB 만 되돌리세요. 토픽까지 지우면 재현이 어려워집니다.

mq "TRUNCATE processed_message; TRUNCATE outbox_event; DELETE FROM orders;
    UPDATE inventory SET quantity = 1000;"

13-1. 지금까지의 결론 — 정확히 한 번은 "설정"이 아니라 "설계"다

Step 07 부터 Step 12 까지, 우리는 계속 같은 벽에 부딪혔습니다.

  • Step 06 — 커밋 시점을 바꿔도 처리와 커밋 사이의 틈은 없어지지 않습니다
  • Step 07 — 재시도는 같은 메시지를 여러 번 처리한다는 뜻입니다
  • Step 08 — 논블로킹 재시도는 순서를 포기하는 대신 중복 가능성을 늘립니다
  • Step 09 — KafkaTransactionManager 는 Kafka 안에서만 원자적이고, DB 와는 원자적이지 않습니다
  • Step 12 — 랙이 0 이어도 처리가 성공했다는 뜻은 아닙니다

결론은 하나입니다. Kafka 로 만들 수 있는 현실적인 보장은 at-least-once 입니다. exactly-once-semantics(EOS)는 "Kafka → 처리 → Kafka" 로 닫힌 경로에서만 성립하고, 우리처럼 DB 와 외부 API 가 끼면 성립하지 않습니다.

그래서 실무의 답은 이렇게 뒤집힙니다.

"메시지가 한 번만 오게 만든다"가 아니라, "여러 번 와도 결과가 같게 만든다."

이것이 멱등성(idempotency)입니다. 이 스텝은 그 멱등성을 세 축으로 나눠 각각 구현합니다.

질문해법
① 멱등 컨슈머같은 메시지를 두 번 받으면?처리 이력 테이블 + UNIQUE 제약13-2 ~ 13-4
② 원자적 발행DB 는 커밋됐는데 발행이 실패하면?Transactional Outbox13-5
③ 순서 보장이벤트가 뒤집혀 도착하면?키 설계 = 순서 설계13-6, 13-7

세 축은 독립적이지 않습니다. ② 를 도입하면 중복 발행이 늘어나므로 ① 이 필수가 되고, ① 을 병렬화하면 ③ 이 깨집니다. 순서대로 봅니다.


13-2. 멱등 컨슈머 — 처리 이력 테이블

가장 단순하고 가장 널리 쓰이는 구현입니다. "이 메시지를 처리한 적이 있는가"를 DB 에 기록합니다.

CREATE TABLE processed_message (
  message_id     VARCHAR(64) NOT NULL PRIMARY KEY,
  consumer_group VARCHAR(64) NOT NULL,
  processed_at   DATETIME(3) NOT NULL DEFAULT CURRENT_TIMESTAMP(3)
) ENGINE=InnoDB;

메시지 ID 를 무엇으로 할 것인가

여기서 대부분의 설계가 갈립니다. 후보는 셋입니다.

후보형태장점치명적 문제
topic-partition-offsetorders-1-42추가 작업 0. 항상 존재재발행하면 오프셋이 바뀝니다. DLT 재처리·Outbox 재발행·토픽 마이그레이션에서 같은 사건이 다른 ID 가 되어 중복 처리됩니다
② 비즈니스 키ORD-0007:OrderCreated사람이 읽을 수 있음. 재발행에 강함같은 타입이 여러 번 정당하게 발생하면 못 씁니다(StockAdjusted 를 두 번 보내면 두 번째가 스킵됨). 이벤트 타입 설계에 종속됩니다
프로듀서가 만든 IDUUID 또는 outbox 행 PK사건 하나 = ID 하나. 재발행해도 동일. 토픽·파티션과 무관프로듀서가 헤더에 실어 줘야 합니다(= 규약이 필요)

③ 을 씁니다. 우리는 13-5 에서 Outbox 를 만들 것이므로, outbox_event.id 를 그대로 메시지 ID 로 씁니다. 릴레이가 재발행해도 행 PK 는 안 바뀌므로 ③ 의 조건을 정확히 만족합니다. Outbox 없이 직접 발행할 때는 UUID.randomUUID() 를 헤더에 실으면 됩니다.

public static final String HDR_MESSAGE_ID = "messageId";

ProducerRecord<String, OrderCreated> record =
        new ProducerRecord<>("orders.outbox", null, event.orderId(), event);
record.headers().add(HDR_MESSAGE_ID, ("OBX-" + row.id()).getBytes(StandardCharsets.UTF_8));
kafkaTemplate.send(record);

구현 — INSERT 를 먼저 시도한다

private static final String INSERT_MARK = """
        INSERT INTO processed_message (message_id, consumer_group)
        VALUES (?, ?)
        """;

boolean markProcessed(String messageId, String group) {
    try {
        jdbc.update(INSERT_MARK, group + ":" + messageId, group);
        return true;                     // 처음 보는 메시지
    } catch (DuplicateKeyException e) {
        return false;                    // 이미 처리함 → 스킵
    }
}

SELECT 가 없습니다. INSERT 를 먼저 던지고, 실패하면 중복으로 간주합니다.

결과 (같은 메시지를 두 번 소비)

INFO 14107 --- [ntainer#0-0-C-1] c.e.o.s13.IdempotentInventoryListener    : consume messageId=OBX-7 order=ORD-0007 sku=SKU-002 qty=3
INFO 14107 --- [ntainer#0-0-C-1] c.e.o.s13.IdempotentInventoryListener    : deducted SKU-002 -3 (remaining=997)
INFO 14107 --- [ntainer#0-0-C-1] c.e.o.s13.IdempotentInventoryListener    : consume messageId=OBX-7 order=ORD-0007 sku=SKU-002 qty=3
WARN 14107 --- [ntainer#0-0-C-1] c.e.o.s13.IdempotentInventoryListener    : duplicate messageId=OBX-7 group=s13-inventory — skipped

두 번째는 deducted 로그가 없습니다. 재고도 997 그대로입니다.

⚠️ 함정 — processed_message 의 PK 가 message_id 하나뿐이다 프로젝트 스키마의 PK 는 message_id 단독입니다. 그런데 컨슈머 그룹은 여럿입니다(s13-inventory, s13-notification). 재고 그룹이 OBX-7 을 먼저 처리해 INSERT 해 버리면, 알림 그룹이 그 메시지를 처음 보는데도 DuplicateKeyException 이 나서 스킵됩니다. 팬아웃 컨슈머의 절반이 조용히 아무 일도 안 하게 됩니다. 예외도, 랙도, 아무 신호가 없습니다. 해결 A — 저장 값을 group + ":" + messageId 로 만들어 그룹별로 다른 키가 되게 합니다(위 코드). 해결 B — 스키마를 고칠 수 있다면 이쪽이 정석입니다.

ALTER TABLE processed_message DROP PRIMARY KEY,
  ADD PRIMARY KEY (message_id, consumer_group);

이 코스는 스키마를 고정해야 하므로 해결 A 로 갑니다. 실무에서는 B 를 쓰세요.

SELECT 로 먼저 확인하면 안 되는가

이렇게 쓰고 싶어집니다.

// ❌ 틀린 코드
Integer cnt = jdbc.queryForObject(
        "SELECT COUNT(*) FROM processed_message WHERE message_id = ?", Integer.class, key);
if (cnt != null && cnt > 0) {
    log.warn("duplicate — skipped");
    return;
}
deductStock(event);                                        // ← ①
jdbc.update("INSERT INTO processed_message ...", key, group);  // ← ②

읽기엔 자연스럽지만 SELECTINSERT 사이가 비어 있습니다. 컨슈머 스레드가 3개(= concurrency: 3)이거나, 인스턴스가 2대이거나, 재시도로 같은 메시지가 겹치면 이렇게 됩니다.

스레드 A                          스레드 B
────────────────────────────────────────────────────
SELECT → 0건                      
                                  SELECT → 0건        ← 아직 A 가 INSERT 하기 전
deductStock()  (-3)               
                                  deductStock()  (-3) ← 두 번 차감
INSERT OK                         
                                  INSERT → DuplicateKey (이미 늦음)

결과 (같은 메시지를 두 스레드가 동시에 소비하도록 강제한 재현)

INFO 14107 --- [       pool-2-t-1] c.e.o.s13.RaceDemo                       : [A] select count=0
INFO 14107 --- [       pool-2-t-2] c.e.o.s13.RaceDemo                       : [B] select count=0
INFO 14107 --- [       pool-2-t-1] c.e.o.s13.RaceDemo                       : [A] deducted SKU-002 -3
INFO 14107 --- [       pool-2-t-2] c.e.o.s13.RaceDemo                       : [B] deducted SKU-002 -3
INFO 14107 --- [       pool-2-t-1] c.e.o.s13.RaceDemo                       : [A] insert ok
WARN 14107 --- [       pool-2-t-2] c.e.o.s13.RaceDemo                       : [B] insert duplicate — but stock already deducted twice

mq "SELECT quantity FROM inventory WHERE sku='SKU-002';"994. 997 이어야 하는데 6개가 빠졌습니다. 100회 반복 실험에서 11회 중복 차감이 발생했습니다.

INSERT-first 방식으로 같은 실험을 하면 0회입니다. 판정을 애플리케이션이 아니라 DB 의 UNIQUE 인덱스에 맡겼기 때문입니다. UNIQUE 제약은 원자적이고, 경합 구간이 존재하지 않습니다.

💡 실무 팁 — "확인하고 하기"가 아니라 "하고 나서 실패를 처리하기" SELECT-then-INSERT, exists()-then-save(), containsKey()-then-put() 은 전부 같은 종류의 버그입니다. 동시성 판정은 제약 조건(UNIQUE / PK) 이나 원자적 연산에 위임하고, 애플리케이션은 실패를 받아 처리하는 쪽이 항상 안전합니다.


13-3. ⚠️ 함정: 멱등 처리와 비즈니스 처리가 다른 트랜잭션이면 소용없다

INSERT-first 로 바꿔도, 이력 기록과 실제 처리가 다른 트랜잭션이면 여전히 깨집니다.

// ❌ 틀린 코드 — 트랜잭션이 둘로 쪼개져 있다
@KafkaListener(topics = "orders.outbox", groupId = "s13-inventory")
public void onMessage(OrderCreated e, @Header("messageId") String messageId) {
    if (!markProcessed(messageId, GROUP)) return;   // 트랜잭션 #1 — 커밋됨
    inventoryService.deduct(e.sku(), e.quantity()); // 트랜잭션 #2 — 여기서 죽으면?
}

markProcessed 가 커밋된 직후 프로세스가 죽으면 이렇게 됩니다.

시점processed_messageinventory결과
markProcessed 커밋OBX-7 있음1000 (미차감)
프로세스 killOBX-7 있음1000오프셋 미커밋
재시작 → 같은 메시지 재소비OBX-7 있음 → 스킵1000재고가 영원히 안 빠짐

결과

INFO 14107 --- [ntainer#0-0-C-1] c.e.o.s13.SplitTxDemo                    : marked OBX-7 (tx#1 committed)
ERROR 14107 --- [ntainer#0-0-C-1] c.e.o.s13.SplitTxDemo                   : simulated crash before deduct
...앱 재시작...
INFO 14107 --- [ntainer#0-0-C-1] o.s.k.l.KafkaMessageListenerContainer    : s13-inventory: partitions assigned: [orders.outbox-0, orders.outbox-1, orders.outbox-2]
WARN 14107 --- [ntainer#0-0-C-1] c.e.o.s13.SplitTxDemo                    : duplicate messageId=OBX-7 group=s13-inventory — skipped

DB 를 확인하면 processed_message 는 1행, inventory.SKU-0021000 그대로입니다.

이력은 남았는데 처리는 안 됐습니다. 그리고 이력이 남았으므로 다시는 처리되지 않습니다. 메시지 하나가 조용히 증발한 것입니다. 로그에 ERROR 도 없습니다(재시작 후 로그는 skipped 뿐입니다).

반대 순서(처리 먼저, 이력 나중)도 대칭적으로 깨집니다. 재고는 빠졌는데 이력이 없으니, 재소비 시 두 번 차감됩니다.

해결 — 하나의 트랜잭션으로 묶는다

@KafkaListener(topics = "orders.outbox", groupId = "s13-inventory")
public void onMessage(OrderCreated e, @Header(HDR_MESSAGE_ID) String messageId) {
    inventoryTx.consumeOnce(messageId, GROUP, e);
}

// 별도 빈이어야 합니다 (자기 호출은 프록시를 안 탑니다)
@Transactional                                     // ← 이력 + 비즈니스가 한 트랜잭션
public void consumeOnce(String messageId, String group, OrderCreated e) {
    try {
        jdbc.update(INSERT_MARK, group + ":" + messageId, group);
    } catch (DuplicateKeyException dup) {
        log.warn("duplicate messageId={} group={} — skipped", messageId, group);
        return;                                    // 롤백 아님. 정상 종료
    }
    int updated = jdbc.update(
            "UPDATE inventory SET quantity = quantity - ? WHERE sku = ? AND quantity >= ?",
            e.quantity(), e.sku(), e.quantity());
    if (updated == 0) {
        throw new OutOfStockException(e.sku());    // ← 롤백 → 이력도 함께 사라짐
    }
}

핵심은 마지막 줄입니다. 재고가 부족해 예외가 나면 processed_message INSERT 도 함께 롤백됩니다. 그래서 재시도가 정상적으로 동작하고, 재시도를 다 소진하면 DLT 로 갑니다(Step 07). 만약 이력만 별도 트랜잭션이었다면 첫 시도에서 이력이 남아 재시도가 전부 스킵되고, DLT 에도 안 갑니다.

결과 (재고 부족 케이스)

INFO 14107 --- [ntainer#0-2-C-1] c.e.o.s13.InventoryTx                    : marked OBX-31 (tx started)
WARN 14107 --- [ntainer#0-2-C-1] c.e.o.s13.InventoryTx                    : out of stock SKU-001 need=4 → rollback
INFO 14107 --- [ntainer#0-2-C-1] c.e.o.s13.InventoryTx                    : marked OBX-31 (tx started)     ← 재시도 1
WARN 14107 --- [ntainer#0-2-C-1] c.e.o.s13.InventoryTx                    : out of stock SKU-001 need=4 → rollback   (재시도 2 도 동일)
ERROR 14107 --- [ntainer#0-2-C-1] o.s.k.l.DefaultErrorHandler             : Backoff FixedBackOff{interval=1000, currentAttempts=3, maxAttempts=3} exhausted for orders.outbox-2@31
INFO 14107 --- [ntainer#0-2-C-1] o.s.k.l.DeadLetterPublishingRecoverer    : Successful dead-letter publication: orders.outbox.DLT-2@0

marked 가 3번 찍혔습니다. 매번 INSERT 하고 매번 롤백했다는 뜻입니다. 이것이 정상입니다.

💡 실무 팁 — @Transactional 은 리스너 메서드에 직접 붙여도 됩니다. 단, DefaultErrorHandler 의 재시도는 트랜잭션 에서 돌아야 합니다. 리스너 메서드에 @Transactional 을 붙이면 매 시도마다 새 트랜잭션이 열리고 닫히므로 문제없습니다. 반대로 KafkaTransactionManager 로 감싸면 재시도 전체가 한 Kafka 트랜잭션에 들어가 롤백 범위가 달라집니다(Step 09 참고).


13-4. ⚠️ 함정: 처리 이력 테이블이 무한히 커진다

processed_message삭제하는 코드를 아무도 안 짜면 영원히 자랍니다.

초당 500건 처리 서비스라면 하루 4,320만 행입니다. 한 달이면 13억 행. PK 인덱스가 메모리를 벗어나는 순간, INSERT 한 건마다 디스크 랜덤 I/O 가 생기고 컨슈머 처리량이 붕괴합니다.

행 수PK 인덱스 크기(추정)markProcessed p99컨슈머 처리량
100만46 MB0.8 ms4,200 msg/s
1,000만460 MB1.4 ms3,900 msg/s
1억4.6 GB27 ms310 msg/s

1억 행 구간에서 13배 느려집니다. 에러는 안 납니다. 그냥 느려집니다.

보관 기간을 정하는 기준

보관 기간 ≥ "같은 메시지가 다시 올 수 있는 최대 기간"

그 최대 기간은 대개 토픽 retention 입니다. 오프셋을 최과거로 리셋하는 최악의 경우, retention 안에 있는 모든 메시지가 다시 옵니다. kt --describe --topic orders.outbox --all | grep retention 으로 확인하면 retention.ms=604800000, 즉 7일입니다. 여유를 둬 14일 보관으로 정합니다. Outbox 재발행·DLT 재처리 지연까지 흡수하려면 retention 의 2배가 안전합니다.

해결 A — TTL 배치 삭제

@Scheduled(cron = "0 10 3 * * *")           // 매일 03:10
public void purge() {
    int total = 0, deleted;
    do {
        deleted = jdbc.update(                          // ★ LIMIT 로 쪼갠다
                "DELETE FROM processed_message WHERE processed_at < ? LIMIT 5000",
                Timestamp.from(Instant.now().minus(14, ChronoUnit.DAYS)));
        total += deleted;
    } while (deleted == 5000);
    log.info("purged {} rows from processed_message", total);
}

결과

INFO 14107 --- [   scheduling-1] c.e.o.s13.ProcessedMessagePurger         : purged 4318227 rows from processed_message

한 번에 수천만 행을 지우면 긴 트랜잭션이 undo 로그를 부풀리고 복제 지연을 만듭니다. LIMIT 분할이 핵심입니다.

해결 B — 파티셔닝 후 파티션 DROP

행이 억 단위라면 DELETE 자체가 부담입니다. 날짜로 RANGE 파티셔닝하고 파티션을 통째로 DROP 하면 즉시 끝납니다.

ALTER TABLE processed_message
  DROP PRIMARY KEY,
  ADD PRIMARY KEY (message_id, processed_at),      -- 파티션 키는 PK 에 포함돼야 합니다
  PARTITION BY RANGE (TO_DAYS(processed_at)) (
    PARTITION p20250101 VALUES LESS THAN (TO_DAYS('2025-01-02')),
    PARTITION p20250102 VALUES LESS THAN (TO_DAYS('2025-01-03')),
    PARTITION pmax      VALUES LESS THAN MAXVALUE
  );

ALTER TABLE processed_message DROP PARTITION p20250101;   -- 0.02 sec

DELETE 4백만 행이 41초 걸린 자리에서, DROP PARTITION0.02초입니다.

💡 파티션 키가 PK 에 포함돼야 한다는 제약, 파티션 프루닝, pmax 관리 등 파티셔닝의 전체 그림은 MySQL 코스 Step 21 을 참고하세요.


13-5. Transactional Outbox 패턴

Step 09 에서 미뤄 둔 문제입니다. "DB 커밋과 Kafka 발행을 어떻게 원자적으로 만드는가."

답은 "만들지 않는다" 입니다. Kafka 를 트랜잭션에서 완전히 빼내고, 발행 의도만 DB 에 같이 저장합니다.

┌──────────────────────────────────────────────────────────┐
│  OrderService.createOrder()      @Transactional          │
│   ① INSERT INTO orders        (비즈니스 상태)             │
│   ② INSERT INTO outbox_event  (발행 의도)                 │
│   ↑ 같은 DB, 같은 트랜잭션 → 둘 다 되거나 둘 다 안 되거나  │
│   ↑ Kafka 호출이 여기에 전혀 없다  ★핵심★                 │
└────────────────────────┬─────────────────────────────────┘
                         │ COMMIT →  outbox_event(published_at IS NULL)
                         ▼          @Scheduled(fixedDelay=500)
              ┌──────────────────────┐  SELECT ... FOR UPDATE SKIP LOCKED
              │   OutboxRelay        │  ③ kafkaTemplate.send().join()
              │                      │  ④ UPDATE published_at = NOW(3)
              └──────────┬───────────┘

                 topic: orders.outbox

프로듀서 쪽 — OrderService

@Service
public class OrderService {

    @Transactional                                   // DataSourceTransactionManager
    public String createOrder(OrderCreated e) {
        jdbc.update("INSERT INTO orders (order_id, customer_id, amount, status)"
                        + " VALUES (?, ?, ?, 'CREATED')",
                e.orderId(), e.customerId(), e.amount());

        jdbc.update("INSERT INTO outbox_event (aggregate_id, event_type, payload)"
                        + " VALUES (?, 'OrderCreated', ?)",
                e.orderId(), toJson(e));

        log.info("order {} created (outbox staged)", e.orderId());
        return e.orderId();
    }
}

kafkaTemplate 이 이 클래스에 아예 주입되지 않은 것이 이 패턴의 전부입니다. Kafka 가 죽어 있어도 주문 생성은 성공합니다. 발행은 나중에 릴레이가 책임집니다.

curl -s -XPOST localhost:8080/orders/1     # 릴레이는 정지시켜 둔 상태
mq "SELECT id, aggregate_id, event_type, published_at FROM outbox_event;"

결과 (before)

+----+--------------+--------------+--------------+
| id | aggregate_id | event_type   | published_at |
+----+--------------+--------------+--------------+
|  1 | ORD-0001     | OrderCreated | NULL         |
+----+--------------+--------------+--------------+

릴레이 쪽 — OutboxRelay

@Component
public class OutboxRelay {

    private static final String PICK = """
            SELECT id, aggregate_id, event_type, payload
            FROM outbox_event
            WHERE published_at IS NULL
            ORDER BY id
            LIMIT 100
            FOR UPDATE SKIP LOCKED
            """;

    @Scheduled(fixedDelay = 500)
    @Transactional                                    // SELECT ~ UPDATE 가 한 트랜잭션
    public void relay() {
        List<OutboxRow> rows = jdbc.query(PICK, OUTBOX_ROW_MAPPER);
        if (rows.isEmpty()) return;

        for (OutboxRow row : rows) {
            OrderCreated event = fromJson(row.payload());
            ProducerRecord<String, OrderCreated> rec =
                    new ProducerRecord<>("orders.outbox", null, row.aggregateId(), event);
            rec.headers().add(HDR_MESSAGE_ID, ("OBX-" + row.id()).getBytes(UTF_8));
            rec.headers().add("eventType", row.eventType().getBytes(UTF_8));

            kafkaTemplate.send(rec).join();           // ★ 발행 확인까지 대기
            jdbc.update("UPDATE outbox_event SET published_at = NOW(3) WHERE id = ?", row.id());
        }
        log.info("relayed {} events (id {}..{})",
                rows.size(), rows.get(0).id(), rows.get(rows.size() - 1).id());
    }
}

결과

INFO 14107 --- [   scheduling-1] c.e.o.s13.OutboxRelay                    : relayed 1 events (id 1..1)
INFO 14107 --- [ntainer#0-1-C-1] c.e.o.s13.InventoryTx                    : marked OBX-1 (tx started)
INFO 14107 --- [ntainer#0-1-C-1] c.e.o.s13.InventoryTx                    : deducted SKU-002 -2 (remaining=998)

결과 (after — SELECT id, created_at, published_at FROM outbox_event)

+----+-------------------------+-------------------------+
| id | created_at              | published_at            |
+----+-------------------------+-------------------------+
|  1 | 2025-06-02 14:31:07.204 | 2025-06-02 14:31:07.463 |
+----+-------------------------+-------------------------+

발행 지연 259ms. fixedDelay=500 이므로 평균 250ms, 최대 500ms + 발행 시간입니다. 100건 측정 결과 p50 = 261ms, p99 = 638ms 였습니다.

💡 kafkaTemplate.send(rec).join()join() 을 빼면 안 됩니다. Step 02 에서 봤듯 send 는 예외 없이 리턴하고, 실패는 CompletableFuture 안에만 있습니다. 확인 없이 published_at 을 갱신하면 발행 안 된 이벤트가 발행됨으로 표시됩니다. 릴레이는 배치 처리량보다 정확성이 중요한 자리이므로 join() 의 비용을 감수합니다.

FOR UPDATE SKIP LOCKED — 다중 인스턴스 안전장치

서비스를 2대 띄우면 릴레이도 2개입니다. 둘이 같은 행을 집으면 모든 이벤트가 2번 발행됩니다.

SKIP LOCKED 는 "다른 트랜잭션이 잠근 행은 기다리지 말고 건너뛰라"는 뜻입니다.

방식인스턴스 A인스턴스 B결과
잠금 없음id 1~100 선택id 1~100 선택전량 중복 발행
FOR UPDATE (SKIP 없음)id 1~100 잠금A 가 커밋할 때까지 블로킹중복은 없지만 처리량 절반, 락 대기 타임아웃 위험
FOR UPDATE SKIP LOCKEDid 1~100 잠금id 101~200 선택중복 없이 병렬

결과 (인스턴스 2대, SKIP LOCKED 있음 / --server.port 만 다르게 기동)

[instance-A] INFO 14107 --- [   scheduling-1] c.e.o.s13.OutboxRelay : relayed 100 events (id 1..100)
[instance-B] INFO 14231 --- [   scheduling-1] c.e.o.s13.OutboxRelay : relayed 100 events (id 101..200)
[instance-A] INFO 14107 --- [   scheduling-1] c.e.o.s13.OutboxRelay : relayed 100 events (id 201..300)

SKIP LOCKED 를 빼고 같은 실험을 한 뒤 kcc --topic orders.outbox --from-beginning | wc -l 를 세면 573 이 나옵니다. 300건이 573건이 됐습니다. 273건이 중복입니다.

⚠️ 함정 — Outbox 는 exactly-once 가 아니라 at-least-once 다 kafkaTemplate.send().join() 이 성공한 직후, UPDATE published_at 을 하기 전에 프로세스가 죽으면? 그 이벤트는 published_at IS NULL 로 남아 있으므로 재시작 후 다시 발행됩니다. 이건 버그가 아니라 이 패턴의 정의된 동작입니다. 발행과 갱신을 원자적으로 만들 방법은 없습니다(그게 가능했다면 애초에 Outbox 가 필요 없었습니다). 그래서 13-2 의 멱등 컨슈머가 선택이 아니라 필수입니다. Outbox 를 도입하면서 멱등 컨슈머를 안 만들면, 중복을 줄이려고 도입한 패턴이 중복을 늘리는 결과가 됩니다. 재현: 릴레이의 send().join()UPDATE 사이에 if (row.id() == 7) System.exit(1); 을 넣고 재시작해 보세요.

폴링 릴레이 vs CDC(Debezium)

항목폴링 릴레이 (이 스텝)CDC — Debezium
발행 지연p50 261ms / p99 638ms (fixedDelay=500)p50 12ms / p99 34ms (binlog tail)
DB 부하유휴 시에도 초당 2회 인덱스 스캔. 인스턴스 수만큼 배수애플리케이션 쿼리 0. binlog 읽기만
published_at 갱신필요 (쓰기 부하 + 테이블 청소 필요)불필요. outbox 행을 INSERT 직후 DELETE 해도 됨
순서 보장ORDER BY id인스턴스 내에서만 보장. SKIP LOCKED 로 인스턴스 간 순서는 뒤섞임binlog 순서 = 커밋 순서. 전역 순서 보장
운영 복잡도낮음. 코드 40줄, 추가 인프라 0높음. Kafka Connect 클러스터 + binlog 권한 + 스키마 관리
장애 지점애플리케이션과 운명 공동체별도 컴포넌트가 하나 늘어남
적합한 규모~수천 TPS수만 TPS 이상, 또는 순서가 엄격한 도메인

💡 처음엔 폴링으로 시작하세요. 코드 40줄로 원자성 문제가 해결됩니다. CDC 는 지연이 문제가 되거나 릴레이의 DB 부하가 눈에 띌 때 옮겨 가면 됩니다. 컨슈머 쪽 코드는 하나도 안 바뀝니다. Kafka Connect 와 Debezium 커넥터 설정은 Kafka 코스 Step 12 에서 다룹니다.


13-6. 순서 보장 전략 — 키 설계가 곧 순서 설계

Kafka 가 보장하는 순서는 딱 하나입니다.

하나의 파티션 안에서, 프로듀서가 보낸 순서대로 컨슈머가 읽는다.

그 외에는 아무것도 보장되지 않습니다. 그리고 이 하나뿐인 보장마저 다섯 가지 방법으로 깨집니다.

#경로무엇이 깨지나방어
키가 없거나 잘못됨같은 주문 이벤트가 다른 파티션으로 분산 → 순서 무의미키 = 순서를 지켜야 하는 단위(= orderId). 절대 null 로 두지 않기
파티션 수 변경hash(key) % partitions 가 통째로 바뀜. 과거 이벤트는 옛 파티션, 신규는 새 파티션파티션 수는 처음에 넉넉히. 늘려야 하면 새 토픽 + 마이그레이션
max.in.flight > 1 + 재시도배치 1 실패 → 재전송하는 사이 배치 2 가 먼저 기록됨enable.idempotence=true (시퀀스 번호로 브로커가 재정렬. max.in.flight<=5 까지 안전)
논블로킹 재시도 (Step 08)실패 메시지가 retry 토픽을 돌아 한참 뒤에 도착순서가 필요하면 @RetryableTopic 대신 DefaultErrorHandler 블로킹 재시도
컨슈머 쪽 병렬 처리파티션에서는 순서대로 읽었는데 워커 스레드가 뒤바꿈13-7 참고. 키 단위 파티셔닝된 워커 풀

③ 을 확인해 보기

enable.idempotence=false + retries=3 + max.in.flight=5 로 두고, 브로커에 일시적 오류를 주입하면 이렇게 됩니다.

결과 (같은 키 ORD-0007 로 v1 → v2 → v3 순서로 발행)

WARN 14107 --- [ad | producer-1] o.a.k.c.p.internals.Sender               : [Producer clientId=producer-1] Got error produce response with correlation id 12 on topic-partition orders.outbox-1, retrying (2 attempts left). Error: NOT_ENOUGH_REPLICAS
INFO 14107 --- [ntainer#0-1-C-1] c.e.o.s13.OrderingProbe                  : received ORD-0007 version=v2 offset=88
INFO 14107 --- [ntainer#0-1-C-1] c.e.o.s13.OrderingProbe                  : received ORD-0007 version=v3 offset=89
INFO 14107 --- [ntainer#0-1-C-1] c.e.o.s13.OrderingProbe                  : received ORD-0007 version=v1 offset=90

v1 이 맨 뒤에 도착했습니다. 상태를 v1 로 되돌리는 이벤트가 마지막에 적용되면 데이터가 과거로 돌아갑니다. 예외는 없습니다.

enable.idempotence=true (프로젝트 기본값)로 바꾸면 같은 재전송 경고가 뜨는데도 v1 → v2 → v3 (offset 88, 89, 90) 순서로 도착합니다. 브로커가 시퀀스 번호로 재정렬해 줍니다. enable.idempotence=true 는 중복 제거뿐 아니라 순서 보장 장치이기도 합니다.

전역 순서가 필요하면?

방법은 하나뿐입니다. 파티션 1개.

파티션 수순서 범위최대 컨슈머 수실측 처리량
1전역18,100 msg/s
3키 단위323,400 msg/s
12키 단위1291,000 msg/s

전역 순서는 처리량 상한을 컨슈머 1개로 못 박는 것과 같은 말입니다. 거의 항상 잘못된 요구사항이고, 실제로 필요한 건 "같은 주문 안에서의 순서"입니다.

💡 실무 팁 — "순서가 필요하다"는 요구를 받으면 먼저 범위를 물으세요. "무엇과 무엇 사이의 순서인가?" 대답이 orderId 면 키를 orderId 로, customerIdcustomerId 로 잡으면 끝납니다. 키를 정하는 것이 곧 "이 단위 안에서는 순서를 지키고, 이 단위끼리는 병렬 처리하겠다" 는 선언입니다. 파티션 배정 알고리즘 자체는 Kafka 코스 Step 03 을 참고하세요.


13-7. ⚠️ 함정: 리스너 안에서 스레드 풀로 넘기면 순서와 커밋이 동시에 깨진다

"컨슈머가 느리니 병렬로 돌리자"는 발상은 대개 이렇게 구현됩니다.

// ❌ 절대 하면 안 되는 코드
@KafkaListener(topics = "orders.outbox", groupId = "s13-inventory")
public void onMessage(OrderCreated e) {
    executor.submit(() -> inventoryTx.consumeOnce(...));   // 즉시 리턴
}

@Async 를 붙여도 완전히 같습니다. 리스너 메서드가 처리를 시작만 하고 즉시 리턴합니다. 두 가지가 동시에 무너집니다.

① 커밋이 처리보다 먼저 일어납니다. Spring 은 리스너 메서드가 정상 리턴하면 "처리 성공"으로 간주하고 오프셋을 커밋합니다. 워커 스레드가 그 뒤에 실패하면 아무도 모릅니다. 재시도도, DLT 도, DefaultErrorHandler 도 전부 무력화됩니다.

② 순서가 깨집니다. 파티션에서는 순서대로 꺼냈지만 워커 스레드 스케줄링은 순서를 보장하지 않습니다.

결과 (같은 키 ORD-0007 의 v1/v2/v3 + 워커에서 예외)

INFO 14107 --- [ntainer#0-1-C-1] c.e.o.s13.AsyncBadListener               : submitted ORD-0007 v1 (offset=88)
INFO 14107 --- [ntainer#0-1-C-1] c.e.o.s13.AsyncBadListener               : submitted ORD-0007 v2 (offset=89)
INFO 14107 --- [ntainer#0-1-C-1] c.e.o.s13.AsyncBadListener               : submitted ORD-0007 v3 (offset=90)
DEBUG 14107 --- [ntainer#0-1-C-1] o.s.k.l.KafkaMessageListenerContainer   : Committing: {orders.outbox-1=OffsetAndMetadata{offset=91}}
INFO 14107 --- [    worker-pool-3] c.e.o.s13.AsyncBadListener              : processed ORD-0007 v3
INFO 14107 --- [    worker-pool-1] c.e.o.s13.AsyncBadListener              : processed ORD-0007 v1
ERROR 14107 --- [    worker-pool-2] c.e.o.s13.AsyncBadListener             : failed ORD-0007 v2
java.lang.RuntimeException: downstream timeout
	at com.example.order.step13.Practice$AsyncBadListener.lambda$onMessage$0(Practice.java:412)

Committing: offset=91처리 로그보다 먼저 찍혔습니다. 그리고 v2 는 실패했는데 이미 커밋됐으므로 영원히 사라졌습니다. 처리 순서는 v3 → v1 → v2 입니다.

kcg --describe --group s13-inventory

결과

GROUP           TOPIC          PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG  CONSUMER-ID                    HOST         CLIENT-ID
s13-inventory   orders.outbox  0          30              30              0    consumer-s13-inventory-1-e41a  /172.19.0.1  consumer-s13-inventory-1
s13-inventory   orders.outbox  1          91              91              0    consumer-s13-inventory-2-77bd  /172.19.0.1  consumer-s13-inventory-2
s13-inventory   orders.outbox  2          28              28              0    consumer-s13-inventory-3-9c03  /172.19.0.1  consumer-s13-inventory-3

LAG 이 전부 0 입니다. 지표상으로는 완벽히 건강한 서비스입니다. Step 12 의 랙 알림도 울리지 않습니다. 메시지만 사라졌습니다. 이것이 이 코스가 계속 말해 온 "에러 없이 조용히 잃는" 의 최종 형태입니다.

그래도 병렬이 필요하다면

해결 1 — 파티션을 늘리고 concurrency 를 올린다. 대부분 이걸로 끝납니다. Spring 이 파티션 단위로 스레드를 배정하므로 순서도 커밋도 안전합니다.

spring.kafka.listener.concurrency: 3   # 파티션 수 이하로. 초과분은 그냥 놉니다 (Step 03)

해결 2 — 키 단위로 파티셔닝된 워커 풀 + 완료 대기. 파티션을 더 못 늘리는데 처리 시간이 긴 경우입니다.

@KafkaListener(topics = "orders.outbox", groupId = "s13-inventory", batch = "true")
public void onBatch(List<ConsumerRecord<String, OrderCreated>> records, Acknowledgment ack) {
    // 같은 키는 항상 같은 워커로 → 키 안에서 순서 유지
    Map<Integer, List<ConsumerRecord<String, OrderCreated>>> lanes = records.stream()
            .collect(Collectors.groupingBy(r -> Math.abs(r.key().hashCode()) % WORKERS));

    List<CompletableFuture<Void>> futures = lanes.values().stream()
            .map(lane -> CompletableFuture.runAsync(() -> lane.forEach(this::handleOne), executor))
            .toList();

    CompletableFuture.allOf(futures.toArray(CompletableFuture[]::new)).join();  // ★ 반드시 대기
    ack.acknowledge();
}

핵심은 두 줄입니다. groupingBy(key) 로 같은 키를 같은 레인에 묶고, join() 으로 전부 끝날 때까지 리스너가 리턴하지 않습니다. join() 이 없으면 위의 ❌ 코드와 똑같아집니다.

결과

INFO 14107 --- [ntainer#0-1-C-1] c.e.o.s13.LaneWorkerListener             : batch=120 lanes=4 elapsed=310ms (serial estimate 1180ms)
DEBUG 14107 --- [ntainer#0-1-C-1] o.s.k.l.KafkaMessageListenerContainer   : Committing: {orders.outbox-1=OffsetAndMetadata{offset=211}}

1,180ms → 310ms. 3.8배 빨라졌고, 커밋은 처리 뒤에 일어났으며, 키 단위 순서는 유지됩니다.

⚠️ max.poll.interval.ms(기본 5분)를 넘기지 않도록 배치 크기를 조절하세요. 레인 하나가 느리면 배치 전체가 대기합니다. max-poll-records 를 500 에서 100 정도로 낮추는 것이 안전합니다.


13-8. 종합 실습 — 주문 서비스 이벤트 연동

지금까지의 13개 스텝을 하나로 조립합니다.

요구사항

  1. POST /orders/{seq} 로 주문을 생성한다. Kafka 가 죽어 있어도 주문 생성은 성공해야 한다.
  2. 주문 생성은 Outbox 로 OrderCreated 를 발행한다 (13-5)
  3. 재고 컨슈머(s13-inventory)는 멱등하게 재고를 차감한다 (13-2, 13-3)
  4. 재고 부족은 1초 간격 3회 재시도 후 orders.outbox.DLT 로 보낸다 (Step 07)
  5. 알림 컨슈머(s13-notification)는 별도 그룹으로 같은 메시지를 받는다 (팬아웃, Step 03)
  6. 컨슈머 랙과 처리 시간을 Actuator 로 노출한다 (Step 12)

아키텍처

   POST /orders/7

┌───────────────────────────────────────────┐
│ OrderController → OrderService            │
│   @Transactional {                        │   ← Kafka 의존 없음
│     INSERT orders                         │
│     INSERT outbox_event(published_at=NULL)│
│   }                                       │
└───────────────────┬───────────────────────┘  MySQL orderdb

        ┌───────────────────────┐  @Scheduled(fixedDelay=500)
        │  OutboxRelay          │  SELECT ... FOR UPDATE SKIP LOCKED
        └───────────┬───────────┘
                    │ send(key=orderId, header messageId=OBX-{id})

     ┌──────────────────────────────────────────┐
     │ topic: orders.outbox  (partitions=3)     │
     └──────┬──────────────────────┬────────────┘
            │ group=s13-inventory  │ group=s13-notification
            ▼                      ▼
  ┌───────────────────────┐  ┌──────────────────────┐
  │ InventoryListener     │  │ NotificationListener │
  │  @Transactional {     │  │  (멱등 키를 그룹명   │
  │    INSERT processed   │  │   접두사로 분리)     │
  │    UPDATE inventory   │  └──────────────────────┘
  │  }                    │
  └───────────┬───────────┘
              │ OutOfStockException × 3 (FixedBackOff 1000ms)

  ┌───────────────────────┐        ┌───────────────────────────┐
  │ orders.outbox.DLT     │───────▶│ DltInspector / 재처리 도구 │
  └───────────────────────┘        └───────────────────────────┘

  관측: /actuator/metrics/kafka.consumer.fetch.manager.records.lag.max
        /actuator/metrics/order.inventory.process   (Timer)
        outbox.backlog   (커스텀 게이지 — 릴레이 장애의 유일한 신호)

구현 체크리스트

#항목근거 스텝
1orders.outbox(3), orders.outbox.DLT(3) 토픽을 NewTopic 빈으로 선언Step 01, 07
2프로듀서 acks=all + enable.idempotence=true + max.in.flight<=5Step 02, 13-6
3OrderService.createOrder()@Transactional 안에 orders + outbox_event INSERT, Kafka 호출 금지13-5
4OutboxRelaySKIP LOCKED 폴링 + send().join() + published_at 갱신13-5
5발행 시 key = orderId, 헤더 messageId = OBX-{id}13-2, 13-6
6컨슈머 ErrorHandlingDeserializer + trusted.packagesStep 04
7InventoryTx.consumeOnce() — 이력 INSERT 와 재고 UPDATE 를 한 트랜잭션으로13-3
8DefaultErrorHandler(DeadLetterPublishingRecoverer, FixedBackOff(1000, 3))Step 07
9NotificationListener다른 groupId 로 등록 (팬아웃)Step 03
10@Timed / MeterRegistry 로 처리 시간, /actuator/metrics 로 랙 노출Step 12

검증 시나리오

① 정상 흐름

mq "TRUNCATE processed_message; TRUNCATE outbox_event; DELETE FROM orders;
    UPDATE inventory SET quantity = 1000;"

for i in $(seq 1 10); do curl -s -XPOST localhost:8080/orders/$i; echo; done

결과 (애플리케이션 콘솔)

INFO 14107 --- [nio-8080-exec-3] c.e.o.s13.OrderService                   : order ORD-0001 created (outbox staged)
INFO 14107 --- [   scheduling-1] c.e.o.s13.OutboxRelay                    : relayed 10 events (id 1..10)
INFO 14107 --- [ntainer#0-1-C-1] c.e.o.s13.InventoryTx                    : deducted SKU-002 -2 (remaining=998)
INFO 14107 --- [ntainer#0-0-C-1] c.e.o.s13.InventoryTx                    : deducted SKU-003 -3 (remaining=997)
INFO 14107 --- [ntainer#1-1-C-1] c.e.o.s13.NotificationListener           : notify customer=1001 order=ORD-0001
mq "SELECT sku, quantity FROM inventory ORDER BY sku;
    SELECT COUNT(*) AS processed FROM processed_message;
    SELECT COUNT(*) AS unpublished FROM outbox_event WHERE published_at IS NULL;"

결과

+---------+----------+     processed   = 20   ← 재고 10 + 알림 10 (그룹 접두사로 분리)
| sku     | quantity |     unpublished =  0   ← 릴레이가 전부 발행함
+---------+----------+
| SKU-001 |      989 |
| SKU-002 |      989 |
| SKU-003 |      992 |
+---------+----------+

총 차감 30개(989+989+992 = 2970 = 3000 − 30). 기대값과 일치합니다.

② 중복 발행 시 멱등 확인

릴레이가 이미 발행한 이벤트를 강제로 재발행합니다.

mq "UPDATE outbox_event SET published_at = NULL;"     # 전부 미발행으로 되돌림

결과

INFO 14107 --- [   scheduling-1] c.e.o.s13.OutboxRelay                    : relayed 10 events (id 1..10)
WARN 14107 --- [ntainer#0-1-C-1] c.e.o.s13.InventoryTx                    : duplicate messageId=OBX-1 group=s13-inventory — skipped
WARN 14107 --- [ntainer#1-1-C-1] c.e.o.s13.NotificationListener           : duplicate messageId=OBX-1 group=s13-notification — skipped

mq "SELECT sku, quantity FROM inventory ORDER BY sku;"989 / 989 / 992. 변화 없음. 이것이 이 스텝의 목표입니다. 메시지는 20건이 흘렀지만 재고는 한 번만 차감됐습니다. 멱등 컨슈머를 끄고(app.idempotent=false) 같은 실험을 하면 979 / 979 / 984 가 됩니다.

③ 재고 부족으로 DLT

mq "UPDATE inventory SET quantity = 1 WHERE sku = 'SKU-001';"
curl -s -XPOST localhost:8080/orders/3        # ORD-0003, SKU-001, qty=4

결과

INFO 14107 --- [   scheduling-1] c.e.o.s13.OutboxRelay                    : relayed 1 events (id 11..11)
WARN 14107 --- [ntainer#0-2-C-1] c.e.o.s13.InventoryTx                    : out of stock SKU-001 need=4 have=1 → rollback   (× 3)
ERROR 14107 --- [ntainer#0-2-C-1] o.s.k.l.DefaultErrorHandler             : Backoff FixedBackOff{interval=1000, currentAttempts=3, maxAttempts=3} exhausted for orders.outbox-2@11
INFO 14107 --- [ntainer#0-2-C-1] o.s.k.l.DeadLetterPublishingRecoverer    : Successful dead-letter publication: orders.outbox.DLT-2@0

결과 (SELECT message_id FROM processed_message WHERE message_id LIKE '%OBX-11')

+-------------------------+
| message_id              |
+-------------------------+
| s13-notification:OBX-11 |
+-------------------------+

재고 그룹의 이력이 없습니다. 트랜잭션이 롤백됐기 때문입니다(13-3). 재고를 보충하고 DLT 에서 재처리하면 정상 처리됩니다. 알림 그룹은 재고와 무관하므로 정상 처리됐습니다.

④ 릴레이 중단 후 재개

릴레이만 끄고 주문을 쌓습니다.

curl -s -XPOST localhost:8080/admin/relay/stop
for i in $(seq 21 25); do curl -s -XPOST localhost:8080/orders/$i; echo; done
mq "SELECT COUNT(*) AS unpublished FROM outbox_event WHERE published_at IS NULL;"

결과

INFO 14107 --- [nio-8080-exec-5] c.e.o.s13.RelayAdmin                     : relay STOPPED
INFO 14107 --- [nio-8080-exec-6] c.e.o.s13.OrderService                   : order ORD-0021 created (outbox staged)
...
+-------------+
| unpublished |
+-------------+
|           5 |
+-------------+

주문 생성은 전부 성공했습니다. 발행 경로가 죽어도 API 는 200 을 돌려줍니다. 이것이 Outbox 의 값어치입니다.

kcg --describe --group s13-inventory 를 찍으면 세 파티션 모두 LAG 0 입니다. 아직 발행 자체가 안 됐기 때문입니다.

💡 여기서 랙은 0 입니다. 밀린 일은 5건인데 Kafka 랙 지표로는 안 보입니다. Outbox 를 쓰면 SELECT COUNT(*) FROM outbox_event WHERE published_at IS NULL 자체가 지표가 되어야 합니다. Step 12 의 MeterRegistry.gauge("outbox.backlog", ...) 로 노출하세요. 이걸 빼먹으면 릴레이 장애를 아무도 모릅니다.

curl -s -XPOST localhost:8080/admin/relay/start

결과

INFO 14107 --- [nio-8080-exec-7] c.e.o.s13.RelayAdmin                     : relay STARTED
INFO 14107 --- [   scheduling-1] c.e.o.s13.OutboxRelay                    : relayed 5 events (id 11..15)
INFO 14107 --- [ntainer#0-0-C-1] c.e.o.s13.InventoryTx                    : deducted SKU-001 -1 (remaining=988)
...

id 11..15중단 중에 쌓인 것부터 순서대로 발행됐습니다. ORDER BY id 덕분입니다.


13-9. 배포 전 체크리스트 20항목

각 항목은 "이 값이면 무슨 일이 일어나는가" 로 읽으세요. 왼쪽이 위험한 값, 오른쪽이 권장입니다.

#설정 / 코드이 값이면 벌어지는 일권장스텝
1send() 결과를 안 봄발행 실패를 영원히 모름whenComplete 또는 릴레이는 join()02
2매 건 send().get()배치 무력화. 18,400 → 920 msg/s비동기 콜백02
3concurrency > partitions남는 스레드가 경고 없이 놈concurrency ≤ partitions03
4컨슈머에 ErrorHandlingDeserializer 없음포이즌 필 1건이 파티션을 영구 정지항상 감싸기04
5spring.json.add.type.headers=true 유지컨슈머 패키지 리팩터링 순간 전량 실패false + default.type 지정04
6trusted.packages: "*"역직렬화 가젯 공격 표면도메인 패키지만 명시04
7ack-mode: BATCH + 배치 중간 실패성공분 재처리 또는 실패분 커밋RECORD 또는 MANUAL + 멱등06
8enable-auto-commit: true처리 전 커밋 → 조용한 유실false (Spring 이 관리)06
9백오프 총합 > max.poll.interval.ms재시도 중 그룹에서 축출 + 리밸런스 루프총합 < 5분, 길면 논블로킹07
10DefaultErrorHandler 기본값FixedBackOff(0, 9) 로 10회 후 조용히 버림Recoverer 로 DLT 지정07
11DLT 파티션 수 ≠ 원본원본 파티션 정보로 보낼 곳이 없어 실패원본과 동일하게07
12@RetryableTopic 을 순서 민감 토픽에재시도 메시지가 한참 뒤에 도착순서 필요하면 블로킹 재시도08
13@Transactional 로 DB+Kafka 를 묶음원자적이지 않음. 한쪽만 반영Outbox09, 13-5
14read_uncommitted(기본)롤백된 트랜잭션 메시지를 소비isolation.level=read_committed09
15enable.idempotence=false + retries>0재전송 중 순서 역전true (프로젝트 기본값)13-6
16운영 중 파티션 수 변경키→파티션 매핑이 바뀌어 순서 파괴처음에 넉넉히. 바꾸려면 새 토픽03, 13-6
17리스너에서 @Async / executor.submit처리 전 커밋 + 순서 파괴 + LAG 은 0concurrency 또는 레인 워커 + join()13-7
18SELECTINSERT 로 멱등 체크경합 구간에서 중복 처리INSERT-first + DuplicateKeyException13-2
19이력 기록과 비즈니스가 다른 트랜잭션이력만 남고 처리 안 됨 → 영구 스킵한 트랜잭션13-3
20processed_message 정리 없음1억 행에서 처리량 13배 하락TTL 배치 또는 파티션 DROP13-4

추가로 관측 3종을 반드시 노출하세요 (Step 12).

지표왜 필요한가
kafka.consumer.fetch.manager.records.lag.max컨슈머가 밀리는 것
outbox.backlog (커스텀 게이지)릴레이가 죽은 것 — Kafka 랙으로는 안 보임
*.DLT 토픽의 LOG-END-OFFSET 증가율조용히 버려지는 메시지

정리

개념핵심
현실의 보장at-least-once. exactly-once 는 멱등성으로 만드는
멱등 컨슈머processed_message + UNIQUE 제약. INSERT-first, DuplicateKeyException 이면 스킵
메시지 IDtopic-partition-offset ❌ / 비즈니스 키 △ / 프로듀서 생성 ID(outbox PK, UUID)
SELECT-then-INSERT경합 구간이 존재. 100회 중 11회 중복 처리. 판정은 DB 제약에 위임
이력 + 비즈니스반드시 한 트랜잭션. 쪼개면 이력만 남고 처리는 안 된 채 영구 스킵
이력 테이블 청소보관 기간 ≥ 토픽 retention. TTL 배치(LIMIT 분할) 또는 파티션 DROP
Transactional Outbox비즈니스 INSERT + outbox INSERT 를 한 트랜잭션. Kafka 를 트랜잭션에서 배제
OutboxRelayFOR UPDATE SKIP LOCKED + send().join() + published_at 갱신
Outbox 의 성질at-least-once. 갱신 전 크래시 = 재발행. 멱등 컨슈머와 반드시 짝
폴링 vs CDC폴링 p50 261ms / 코드 40줄. CDC p50 12ms / 인프라 1개 추가. 폴링부터 시작
순서 보장 범위파티션 내부만. 전역 순서 = 파티션 1개 = 처리량 8,100 msg/s 상한
키 설계키를 정하는 것이 "이 단위 안에서 순서를 지키겠다"는 선언
순서를 깨는 5가지키 없음 / 파티션 수 변경 / 비멱등 재전송 / 논블로킹 재시도 / 컨슈머 병렬화
@Async 리스너처리 전 커밋 + 순서 파괴. 그런데 LAG 은 0. 최악의 침묵
안전한 병렬화concurrency 우선. 부족하면 키 단위 레인 + allOf(...).join()
Outbox 백로그 지표Kafka 랙으로 안 보임. outbox.backlog 게이지를 별도로 노출

연습문제

Exercise.java 에 6문제가 있습니다. 정답은 Solution.java. 전부 MySQL 결과를 눈으로 확인해야 하는 문제입니다.

  1. messageId 헤더 기반 멱등 컨슈머를 구현하고, 같은 이벤트 10건을 두 번 발행해도 inventory 가 한 번만 차감되는지 검증하기
  2. SELECT-then-INSERT 방식으로 바꾼 뒤 두 스레드로 같은 메시지를 100회 동시 처리해 중복 차감 횟수를 세기
  3. OrderService + OutboxRelay 를 구현하고, Kafka 컨테이너를 내린 상태에서 주문 생성이 성공하는지 확인하기
  4. SKIP LOCKED 를 뺀 릴레이를 두 인스턴스로 띄워 중복 발행 건수를 측정하고, 넣었을 때와 비교하기
  5. 주어진 5개 코드 조각 중 순서 보장이 깨지는 것을 모두 고르고 각각의 이유를 한 줄로 적기
  6. orders.outbox.DLT 를 읽어 원본 토픽으로 재발행하는 재처리 도구를 만들기 (messageId 헤더를 보존해야 멱등성이 유지됩니다)

코스를 마치며

13개 스텝, 하나의 주제였습니다. 에러 없이 조용히 잘못 동작하는 코드를 찾아내는 것.

각 스텝에서 잡은 침묵을 한 줄씩 되짚습니다.

Step조용한 실패신호
01자동 설정이 만든 빈을 모른 채 덮어씀없음
02send() 가 리턴했는데 브로커에 안 감CompletableFuture 안에만
03concurrency 초과분 스레드가 그냥 놈없음
04__TypeId__ 패키지 불일치로 전량 실패파티션 정지
05헤더가 없어 추적 ID 가 끊김없음
06처리 전에 커밋되어 메시지 유실LAG 은 0
07FixedBackOff(0,9) 소진 후 조용히 버림ERROR 한 줄뿐
08재시도 메시지가 순서를 벗어나 도착없음
09DB 는 커밋, Kafka 는 롤백없음
10autoStartup=false 컨테이너를 아무도 안 켬랙 증가만
11Thread.sleep 테스트가 CI 에서만 실패간헐적
12지표는 있는데 아무도 안 봄
13이력만 남고 처리 안 됨 / @Async 로 처리 전 커밋LAG 은 0

공통점이 보이시나요. 거의 모든 항목의 "신호" 칸이 비어 있습니다. Kafka 를 쓰는 일의 절반은 이 빈칸을 채워 넣는 일입니다. 로그를 교재와 대조하라고 반복해서 말한 이유가 이것입니다.

더 공부할 것

주제무엇인가어디서
Kafka Streams토픽 → 토픽 변환을 상태 저장소와 함께. 조인·윈도우 집계·exactly_once_v2Kafka 코스 Step 13
Kafka Connect / CDC코드 없이 DB ↔ Kafka. Debezium 으로 13-5 의 폴링 릴레이 대체Kafka 코스 Step 12
스키마 레지스트리Avro/Protobuf + 호환성 규칙. Step 04 의 __TypeId__ 결합 문제를 계약으로 해결Confluent Schema Registry
Saga 패턴서비스에 걸친 트랜잭션을 보상 트랜잭션의 연쇄로. Outbox 가 그 전제코레오그래피 / 오케스트레이션
이벤트 소싱상태가 아니라 이벤트를 원본으로 저장. 스냅숏·리플레이·CQRS도메인 설계 영역

우선순위를 하나만 고른다면 스키마 레지스트리입니다. 컨슈머가 셋 이상으로 늘어나는 순간, 이벤트 스키마 변경이 가장 자주 터지는 사고 지점이 됩니다.

수고하셨습니다.


실습 파일

이 스텝은 앞의 열두 스텝과 달리 완성된 서비스를 조립합니다. Practice.javastep13-idemstep13-outboxstep13-order 순서로 켜면서 13-2 → 13-5 → 13-8 을 차례로 재현하고, 매 단계마다 MySQL 을 직접 조회해 결과를 확인하세요. 이 스텝은 콘솔 로그만 봐서는 절반밖에 못 봅니다. 그다음 Exercise.java 의 6문제를 풀고 Solution.java 로 대조합니다. 세 파일 모두 com.example.order.step13 패키지에 둡니다.

Practice.java

본문 13-2 ~ 13-8 의 모든 예제를 절 번호 주석과 함께 nested static class 로 담았습니다.

  • 프로필이 5개입니다. step13,step13-idem(13-2 멱등 + 경합 재현), step13,step13-split(13-3 트랜잭션 분리 함정), step13,step13-outbox(13-5 Outbox 전체), step13,step13-async(13-7 @Async 함정), step13,step13-order(13-8 최종 프로젝트). 하나씩만 켜세요. 리스너 그룹이 겹치면 재현이 안 됩니다.
  • Sql 클래스에 이 스텝의 모든 SQL 상수가 모여 있습니다. INSERT_MARK, PICK_OUTBOX(FOR UPDATE SKIP LOCKED), DEDUCT(AND quantity >= ? 조건 포함)를 먼저 읽으면 나머지 코드가 빨리 읽힙니다. 특히 DEDUCTupdated == 0 판정이 재고 부족 검출의 전부입니다.
  • [13-2] RaceDemo일부러 Kafka 를 안 씁니다. 두 스레드를 CountDownLatch 로 동시에 출발시켜 SELECT-then-INSERT 의 경합만 순수하게 재현합니다. ROUNDS = 100 을 돌리고 마지막에 중복 차감 횟수를 집계해 찍습니다. 실행할 때마다 횟수가 달라지는 것 자체가 이 문제의 성질입니다.
  • [13-5] OutboxRelayCRASH_AT_ID 상수를 -1 이 아닌 값(예: 7)으로 바꾸면, 발행 직후 published_at 갱신 전에 System.exit(1) 로 죽습니다. 재시작 후 그 이벤트가 다시 발행되는 것을 확인하는 용도입니다. Outbox 가 at-least-once 라는 것을 몸으로 확인하는 유일한 방법입니다.
  • [13-7] AsyncBadListener 는 잘못된 코드이고 [13-7] LaneWorkerListener 가 올바른 코드입니다. 두 리스너가 같은 프로필에 있지만 그룹이 다르므로 같은 메시지를 각자 받습니다. 커밋 로그를 비교하려면 logging.level.org.springframework.kafka.listener=DEBUG 를 켜세요. Committing: 줄의 위치가 전부입니다.
  • [13-8] OrderControllerRelayAdminspring-boot-starter-web 이 필요합니다. 13-0 에서 추가하지 않았다면 이 프로필은 기동에 실패합니다. RelayAdmin 은 릴레이를 켜고 끄는 플래그만 토글하며, 검증 시나리오 ④ 에서 씁니다.
package com.example.order.step13;

/*
 * ============================================================================
 * Step 13 — 실전 패턴과 최종 프로젝트 : Practice
 * ============================================================================
 *
 * 배치 위치
 *   spring-kafka-lab/src/main/java/com/example/order/step13/Practice.java
 *
 * 사전 준비
 *   1) build.gradle 에 웹 스타터 추가 (13-8 의 REST 엔드포인트에 필요)
 *        implementation 'org.springframework.boot:spring-boot-starter-web'
 *   2) 토픽 생성
 *        kt --create --topic orders.outbox     --partitions 3 --replication-factor 1
 *        kt --create --topic orders.outbox.DLT --partitions 3 --replication-factor 1
 *   3) MySQL 별칭
 *        alias mq='docker exec -i learn-kafka-mysql mysql -ulearner -plearn1234 orderdb -t -e'
 *
 * 실행 (보조 프로필을 하나만 함께 켭니다. 예제끼리 서로 간섭합니다)
 *   ./gradlew bootRun --args='--spring.profiles.active=step13,step13-idem'
 *       → [13-2] 멱등 컨슈머 + SELECT-then-INSERT 경합 재현(RaceDemo)
 *   ./gradlew bootRun --args='--spring.profiles.active=step13,step13-split'
 *       → [13-3] ★함정★ 이력과 비즈니스가 다른 트랜잭션이면 메시지가 영구 스킵된다
 *   ./gradlew bootRun --args='--spring.profiles.active=step13,step13-outbox'
 *       → [13-5] Transactional Outbox 전체 (OrderService + OutboxRelay)
 *   ./gradlew bootRun --args='--spring.profiles.active=step13,step13-async'
 *       → [13-7] ★함정★ @Async 리스너가 처리 전에 커밋한다 vs 레인 워커
 *   ./gradlew bootRun --args='--spring.profiles.active=step13,step13-order'
 *       → [13-8] 최종 프로젝트. REST + Outbox + 멱등 재고 + DLT + 알림 팬아웃
 *
 * 실행 전 DB 초기화 (매번 하세요. 안 하면 이력이 남아 전부 스킵됩니다)
 *   mq "TRUNCATE processed_message; TRUNCATE outbox_event; DELETE FROM orders;
 *       UPDATE inventory SET quantity = 1000;"
 *
 * 실행 중 확인할 CLI
 *   mq  "SELECT sku, quantity FROM inventory ORDER BY sku;"
 *   mq  "SELECT id, aggregate_id, created_at, published_at FROM outbox_event ORDER BY id;"
 *   mq  "SELECT message_id, consumer_group FROM processed_message ORDER BY processed_at;"
 *   kcg --describe --group s13-inventory
 *   kcc --topic orders.outbox --from-beginning --property print.key=true --property print.headers=true
 *
 * 커밋 시점을 눈으로 보려면 (13-7 에 필수)
 *   logging.level.org.springframework.kafka.listener=DEBUG
 * ============================================================================
 */

import com.example.order.domain.OrderCreated;
import com.fasterxml.jackson.databind.ObjectMapper;

import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.ApplicationArguments;
import org.springframework.boot.ApplicationRunner;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Profile;
import org.springframework.dao.DuplicateKeyException;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.jdbc.core.RowMapper;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.core.ConsumerFactory;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.listener.DeadLetterPublishingRecoverer;
import org.springframework.kafka.listener.DefaultErrorHandler;
import org.springframework.kafka.support.Acknowledgment;
import org.springframework.kafka.support.KafkaHeaders;
import org.springframework.messaging.handler.annotation.Header;
import org.springframework.scheduling.annotation.EnableScheduling;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
import org.springframework.transaction.annotation.Transactional;
import org.springframework.util.backoff.FixedBackOff;
import org.springframework.web.bind.annotation.PathVariable;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RestController;

import java.nio.charset.StandardCharsets;
import java.sql.Timestamp;
import java.time.Instant;
import java.time.temporal.ChronoUnit;
import java.util.List;
import java.util.Map;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.stream.Collectors;

public final class Practice {

    private Practice() {
    }

    /* =====================================================================
     * 공통 — 상수, SQL, 예외, 행 매퍼
     * ===================================================================== */

    /** 프로듀서가 만들어 실어 보내는 메시지 식별자 헤더. 13-2 의 후보 ③. */
    public static final String HDR_MESSAGE_ID = "messageId";
    public static final String HDR_EVENT_TYPE = "eventType";

    public static final String TOPIC_OUTBOX = "orders.outbox";
    public static final String TOPIC_OUTBOX_DLT = "orders.outbox.DLT";

    public static final String GROUP_INVENTORY = "s13-inventory";
    public static final String GROUP_NOTIFICATION = "s13-notification";

    /** 이 스텝의 모든 SQL. 먼저 읽고 나머지 코드를 보면 빠릅니다. */
    public static final class Sql {
        private Sql() {
        }

        /** 멱등 판정을 DB 제약에 위임한다. SELECT 가 없다는 점이 핵심. */
        public static final String INSERT_MARK = """
                INSERT INTO processed_message (message_id, consumer_group)
                VALUES (?, ?)
                """;

        /** ⚠️ 13-2 의 틀린 방식. 경합 재현용으로만 씁니다. */
        public static final String COUNT_MARK = """
                SELECT COUNT(*) FROM processed_message WHERE message_id = ?
                """;

        /**
         * quantity >= ? 조건이 재고 부족 검출의 전부입니다.
         * 조건에 안 맞으면 updated == 0 이 되고, 거기서 예외를 던져 트랜잭션을 롤백합니다.
         */
        public static final String DEDUCT = """
                UPDATE inventory SET quantity = quantity - ?
                WHERE sku = ? AND quantity >= ?
                """;

        public static final String INSERT_ORDER = """
                INSERT INTO orders (order_id, customer_id, amount, status)
                VALUES (?, ?, ?, 'CREATED')
                """;

        public static final String INSERT_OUTBOX = """
                INSERT INTO outbox_event (aggregate_id, event_type, payload)
                VALUES (?, ?, ?)
                """;

        /**
         * SKIP LOCKED — 다른 트랜잭션이 잠근 행은 기다리지 말고 건너뛴다.
         * 이것이 없으면 릴레이 인스턴스 2대가 같은 행을 집어 전량 중복 발행합니다.
         */
        public static final String PICK_OUTBOX = """
                SELECT id, aggregate_id, event_type, payload
                FROM outbox_event
                WHERE published_at IS NULL
                ORDER BY id
                LIMIT 100
                FOR UPDATE SKIP LOCKED
                """;

        public static final String MARK_PUBLISHED = """
                UPDATE outbox_event SET published_at = NOW(3) WHERE id = ?
                """;

        public static final String PURGE_MARKS = """
                DELETE FROM processed_message WHERE processed_at < ? LIMIT 5000
                """;

        public static final String OUTBOX_BACKLOG = """
                SELECT COUNT(*) FROM outbox_event WHERE published_at IS NULL
                """;
    }

    /** 재고 부족. 재시도해도 소용없지만, 재고 보충 가능성이 있으므로 재시도 대상으로 둡니다. */
    public static class OutOfStockException extends RuntimeException {
        public OutOfStockException(String sku, int need, int have) {
            super("out of stock %s need=%d have=%d".formatted(sku, need, have));
        }
    }

    public record OutboxRow(long id, String aggregateId, String eventType, String payload) {
    }

    public static final RowMapper<OutboxRow> OUTBOX_ROW_MAPPER = (rs, n) -> new OutboxRow(
            rs.getLong("id"),
            rs.getString("aggregate_id"),
            rs.getString("event_type"),
            rs.getString("payload"));

    /**
     * 멱등 키를 만드는 유일한 자리.
     * ⚠️ processed_message 의 PK 가 message_id 단독이므로 그룹명을 접두사로 붙입니다.
     *    안 붙이면 재고 그룹이 먼저 INSERT 한 순간 알림 그룹이 전부 스킵됩니다(13-2 함정).
     */
    public static String markKey(String group, String messageId) {
        return group + ":" + messageId;
    }

    /* =====================================================================
     * [13-2] 멱등 컨슈머 — 처리 이력 테이블
     * ===================================================================== */

    /**
     * ★ 이력 INSERT 와 재고 UPDATE 가 반드시 같은 트랜잭션이어야 합니다(13-3).
     * 그래서 리스너와 분리된 별도 빈으로 둡니다. 자기 호출(this.method())은 프록시를 안 타서
     * @Transactional 이 통째로 무시됩니다.
     */
    @Component
    public static class InventoryTx {

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

        private final JdbcTemplate jdbc;

        public InventoryTx(JdbcTemplate jdbc) {
            this.jdbc = jdbc;
        }

        @Transactional
        public void consumeOnce(String messageId, String group, OrderCreated e) {
            try {
                jdbc.update(Sql.INSERT_MARK, markKey(group, messageId), group);
                log.info("marked {} (tx started)", messageId);
            } catch (DuplicateKeyException dup) {
                // 롤백이 아니라 정상 종료입니다. 이미 처리한 메시지이므로 할 일이 없습니다.
                log.warn("duplicate messageId={} group={} — skipped", messageId, group);
                return;
            }

            int updated = jdbc.update(Sql.DEDUCT, e.quantity(), e.sku(), e.quantity());
            if (updated == 0) {
                Integer have = jdbc.queryForObject(
                        "SELECT quantity FROM inventory WHERE sku = ?", Integer.class, e.sku());
                log.warn("out of stock {} need={} have={} → rollback", e.sku(), e.quantity(), have);
                // ★ 여기서 던지면 위의 INSERT_MARK 도 함께 롤백됩니다.
                //   그래야 재시도가 정상 동작하고, 소진 후 DLT 로 갑니다.
                throw new OutOfStockException(e.sku(), e.quantity(), have == null ? -1 : have);
            }

            Integer remaining = jdbc.queryForObject(
                    "SELECT quantity FROM inventory WHERE sku = ?", Integer.class, e.sku());
            log.info("deducted {} -{} (remaining={})", e.sku(), e.quantity(), remaining);
        }
    }

    @Configuration
    @Profile("step13-idem")
    public static class IdempotentDemo {

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

        private final InventoryTx tx;

        public IdempotentDemo(InventoryTx tx) {
            this.tx = tx;
        }

        // [13-2] 멱등 컨슈머. 리스너는 메시지를 풀어 넘기기만 하고, 트랜잭션은 InventoryTx 가 엽니다.
        @KafkaListener(topics = TOPIC_OUTBOX, groupId = GROUP_INVENTORY)
        public void onMessage(OrderCreated e,
                              @Header(name = HDR_MESSAGE_ID, required = false) String messageId,
                              @Header(KafkaHeaders.RECEIVED_TOPIC) String topic,
                              @Header(KafkaHeaders.OFFSET) long offset) {
            // 헤더가 없으면 멱등성이 성립하지 않습니다. 조용히 넘어가지 말고 실패시킵니다.
            if (messageId == null) {
                throw new IllegalStateException(
                        "missing %s header at %s@%d".formatted(HDR_MESSAGE_ID, topic, offset));
            }
            log.info("consume messageId={} order={} sku={} qty={}",
                    messageId, e.orderId(), e.sku(), e.quantity());
            tx.consumeOnce(messageId, GROUP_INVENTORY, e);
        }
    }

    /**
     * [13-2] ⚠️ SELECT-then-INSERT 의 경합 재현.
     * Kafka 를 전혀 쓰지 않습니다. 두 스레드를 CountDownLatch 로 동시에 출발시켜
     * "확인하고 하기" 패턴의 빈 구간만 순수하게 드러냅니다.
     */
    @Component
    @Profile("step13-idem")
    public static class RaceDemo implements ApplicationRunner {

        private static final Logger log = LoggerFactory.getLogger(RaceDemo.class);
        private static final int ROUNDS = 100;

        private final JdbcTemplate jdbc;

        public RaceDemo(JdbcTemplate jdbc) {
            this.jdbc = jdbc;
        }

        @Override
        public void run(ApplicationArguments args) throws Exception {
            ExecutorService pool = Executors.newFixedThreadPool(2);
            AtomicInteger doubleDeduct = new AtomicInteger();

            for (int round = 1; round <= ROUNDS; round++) {
                String messageId = "RACE-" + round;
                jdbc.update("UPDATE inventory SET quantity = 1000 WHERE sku = 'SKU-002'");

                CountDownLatch start = new CountDownLatch(1);
                CountDownLatch done = new CountDownLatch(2);

                for (String tag : new String[]{"A", "B"}) {
                    pool.submit(() -> {
                        try {
                            start.await();
                            wrongIdempotentDeduct(tag, messageId);
                        } catch (Exception ignored) {
                            // 재현이 목적이므로 예외는 무시합니다
                        } finally {
                            done.countDown();
                        }
                    });
                }
                start.countDown();
                done.await();

                Integer qty = jdbc.queryForObject(
                        "SELECT quantity FROM inventory WHERE sku = 'SKU-002'", Integer.class);
                if (qty != null && qty < 997) {          // 정상이면 997 (1회 차감)
                    doubleDeduct.incrementAndGet();
                }
            }
            pool.shutdown();

            log.warn("SELECT-then-INSERT: {} / {} rounds double-deducted", doubleDeduct.get(), ROUNDS);
            log.warn("→ 실행할 때마다 숫자가 달라집니다. 그 사실 자체가 이 버그의 성질입니다.");
        }

        /** ❌ 틀린 코드. SELECT 와 INSERT 사이가 비어 있습니다. */
        private void wrongIdempotentDeduct(String tag, String messageId) {
            Integer cnt = jdbc.queryForObject(Sql.COUNT_MARK, Integer.class, messageId);
            log.info("[{}] select count={}", tag, cnt);
            if (cnt != null && cnt > 0) {
                log.warn("[{}] duplicate — skipped", tag);
                return;
            }
            jdbc.update(Sql.DEDUCT, 3, "SKU-002", 3);
            log.info("[{}] deducted SKU-002 -3", tag);
            try {
                jdbc.update(Sql.INSERT_MARK, messageId, "race");
                log.info("[{}] insert ok", tag);
            } catch (DuplicateKeyException dup) {
                log.warn("[{}] insert duplicate — but stock already deducted twice", tag);
            }
        }
    }

    /* =====================================================================
     * [13-3] ⚠️ 함정 — 이력과 비즈니스가 다른 트랜잭션이면 소용없다
     * ===================================================================== */

    @Component
    @Profile("step13-split")
    public static class SplitTxDemo {

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

        /** 이 오프셋을 처리할 때 차감 직전에 죽입니다. -1 이면 비활성. */
        private static final long CRASH_AT_OFFSET = 0L;

        private final JdbcTemplate jdbc;

        public SplitTxDemo(JdbcTemplate jdbc) {
            this.jdbc = jdbc;
        }

        @KafkaListener(topics = TOPIC_OUTBOX, groupId = "s13-split")
        public void onMessage(OrderCreated e,
                              @Header(HDR_MESSAGE_ID) String messageId,
                              @Header(KafkaHeaders.OFFSET) long offset) {
            // ❌ 트랜잭션 #1 — 여기서 커밋됩니다
            if (!markInSeparateTx(messageId)) {
                log.warn("duplicate messageId={} group=s13-split — skipped", messageId);
                return;
            }
            log.info("marked {} (tx#1 committed)", messageId);

            if (offset == CRASH_AT_OFFSET) {
                log.error("simulated crash before deduct");
                Runtime.getRuntime().halt(1);            // 셧다운 훅 없이 즉사
            }

            // ❌ 트랜잭션 #2 — 여기 도달 못 하면 이력만 남습니다
            deductInSeparateTx(e);
        }

        @Transactional(propagation = org.springframework.transaction.annotation.Propagation.REQUIRES_NEW)
        public boolean markInSeparateTx(String messageId) {
            try {
                jdbc.update(Sql.INSERT_MARK, markKey("s13-split", messageId), "s13-split");
                return true;
            } catch (DuplicateKeyException dup) {
                return false;
            }
        }

        @Transactional(propagation = org.springframework.transaction.annotation.Propagation.REQUIRES_NEW)
        public void deductInSeparateTx(OrderCreated e) {
            jdbc.update(Sql.DEDUCT, e.quantity(), e.sku(), e.quantity());
            log.info("deducted {} -{}", e.sku(), e.quantity());
        }
    }

    /* =====================================================================
     * [13-4] 처리 이력 테이블 청소
     * ===================================================================== */

    @Component
    @Profile({"step13-idem", "step13-order"})
    public static class ProcessedMessagePurger {

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

        /** 보관 기간 ≥ 토픽 retention(7일). 재발행·DLT 재처리 여유까지 감안해 2배로 잡습니다. */
        private static final int RETENTION_DAYS = 14;

        private final JdbcTemplate jdbc;

        public ProcessedMessagePurger(JdbcTemplate jdbc) {
            this.jdbc = jdbc;
        }

        @Scheduled(cron = "0 10 3 * * *")
        public void purge() {
            Timestamp cutoff = Timestamp.from(
                    Instant.now().minus(RETENTION_DAYS, ChronoUnit.DAYS));
            int total = 0;
            int deleted;
            do {
                // LIMIT 로 쪼갭니다. 한 번에 수천만 행을 지우면 undo 로그가 폭발하고
                // 복제 지연이 생깁니다.
                deleted = jdbc.update(Sql.PURGE_MARKS, cutoff);
                total += deleted;
            } while (deleted == 5000);
            log.info("purged {} rows from processed_message", total);
        }
    }

    /* =====================================================================
     * [13-5] Transactional Outbox
     * ===================================================================== */

    /**
     * ★ 이 클래스에 KafkaTemplate 이 주입되지 않은 것이 패턴의 전부입니다.
     * Kafka 가 죽어 있어도 주문 생성은 성공합니다.
     */
    @Component
    public static class OrderService {

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

        private final JdbcTemplate jdbc;
        private final ObjectMapper mapper;

        public OrderService(JdbcTemplate jdbc, ObjectMapper mapper) {
            this.jdbc = jdbc;
            this.mapper = mapper;
        }

        @Transactional
        public String createOrder(OrderCreated e) {
            jdbc.update(Sql.INSERT_ORDER, e.orderId(), e.customerId(), e.amount());
            jdbc.update(Sql.INSERT_OUTBOX, e.orderId(), "OrderCreated", toJson(e));
            log.info("order {} created (outbox staged)", e.orderId());
            return e.orderId();
        }

        private String toJson(OrderCreated e) {
            try {
                return mapper.writeValueAsString(e);
            } catch (Exception ex) {
                throw new IllegalStateException("cannot serialize " + e.orderId(), ex);
            }
        }
    }

    /** 릴레이를 켜고 끄는 플래그. 검증 시나리오 ④ 에서 씁니다. */
    @Component
    public static class RelayAdmin {

        private static final Logger log = LoggerFactory.getLogger(RelayAdmin.class);
        private final AtomicBoolean running = new AtomicBoolean(true);

        public boolean isRunning() {
            return running.get();
        }

        public void stop() {
            running.set(false);
            log.info("relay STOPPED");
        }

        public void start() {
            running.set(true);
            log.info("relay STARTED");
        }
    }

    @Configuration
    @EnableScheduling
    @Profile({"step13-outbox", "step13-order"})
    public static class OutboxRelay {

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

        /**
         * ★ 실험용. 이 id 를 발행한 직후, published_at 갱신 전에 프로세스를 죽입니다.
         * 재시작하면 그 이벤트가 다시 발행됩니다 = Outbox 는 at-least-once.
         * 평소에는 -1 로 두세요.
         */
        private static final long CRASH_AT_ID = -1L;

        private final JdbcTemplate jdbc;
        private final ObjectMapper mapper;
        private final KafkaTemplate<String, Object> kafkaTemplate;
        private final RelayAdmin admin;

        public OutboxRelay(JdbcTemplate jdbc, ObjectMapper mapper,
                           KafkaTemplate<String, Object> kafkaTemplate, RelayAdmin admin) {
            this.jdbc = jdbc;
            this.mapper = mapper;
            this.kafkaTemplate = kafkaTemplate;
            this.admin = admin;
        }

        @Scheduled(fixedDelay = 500)
        @Transactional                       // SELECT ... FOR UPDATE 부터 UPDATE 까지가 한 트랜잭션
        public void relay() {
            if (!admin.isRunning()) {
                return;
            }
            List<OutboxRow> rows = jdbc.query(Sql.PICK_OUTBOX, OUTBOX_ROW_MAPPER);
            if (rows.isEmpty()) {
                return;
            }

            for (OutboxRow row : rows) {
                OrderCreated event = fromJson(row.payload());

                ProducerRecord<String, Object> rec =
                        new ProducerRecord<>(TOPIC_OUTBOX, null, row.aggregateId(), event);
                // 메시지 ID = outbox 행 PK. 재발행해도 값이 안 바뀌므로 멱등 키로 적합합니다.
                rec.headers().add(HDR_MESSAGE_ID,
                        ("OBX-" + row.id()).getBytes(StandardCharsets.UTF_8));
                rec.headers().add(HDR_EVENT_TYPE,
                        row.eventType().getBytes(StandardCharsets.UTF_8));

                // ★ join() 필수. send 는 예외 없이 리턴하고 실패는 Future 안에만 있습니다(Step 02).
                //   확인 없이 published_at 을 갱신하면 발행 안 된 이벤트가 발행됨으로 표시됩니다.
                kafkaTemplate.send(rec).join();

                if (row.id() == CRASH_AT_ID) {
                    log.error("simulated crash after send, before published_at update (id={})", row.id());
                    Runtime.getRuntime().halt(1);
                }

                jdbc.update(Sql.MARK_PUBLISHED, row.id());
            }
            log.info("relayed {} events (id {}..{})",
                    rows.size(), rows.get(0).id(), rows.get(rows.size() - 1).id());
        }

        /**
         * Kafka 랙으로는 릴레이 장애가 안 보입니다. 이 게이지가 유일한 신호입니다(13-8).
         */
        @Scheduled(fixedDelay = 5000)
        public void reportBacklog() {
            Integer backlog = jdbc.queryForObject(Sql.OUTBOX_BACKLOG, Integer.class);
            if (backlog != null && backlog > 0) {
                log.info("outbox.backlog={}", backlog);
            }
        }

        private OrderCreated fromJson(String payload) {
            try {
                return mapper.readValue(payload, OrderCreated.class);
            } catch (Exception ex) {
                throw new IllegalStateException("cannot deserialize outbox payload", ex);
            }
        }
    }

    /* =====================================================================
     * [13-6] 순서 보장 관측
     * ===================================================================== */

    @Component
    @Profile("step13-outbox")
    public static class OrderingProbe {

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

        @KafkaListener(topics = TOPIC_OUTBOX, groupId = "s13-ordering", concurrency = "1")
        public void onMessage(ConsumerRecord<String, OrderCreated> rec) {
            // 같은 키는 항상 같은 파티션이므로, 여기 도착 순서가 곧 발행 순서여야 합니다.
            log.info("received {} createdAt={} partition={} offset={}",
                    rec.key(), rec.value().createdAt(), rec.partition(), rec.offset());
        }
    }

    /* =====================================================================
     * [13-7] ⚠️ 함정 — @Async / 스레드 풀은 순서와 커밋을 동시에 깨뜨린다
     * ===================================================================== */

    @Component
    @Profile("step13-async")
    public static class AsyncBadListener {

        private static final Logger log = LoggerFactory.getLogger(AsyncBadListener.class);
        private final ExecutorService pool = Executors.newFixedThreadPool(4);

        // ❌ 절대 하면 안 되는 코드
        @KafkaListener(topics = TOPIC_OUTBOX, groupId = "s13-async-bad")
        public void onMessage(ConsumerRecord<String, OrderCreated> rec) {
            log.info("submitted {} (offset={})", rec.key(), rec.offset());
            pool.submit(() -> {
                try {
                    Thread.sleep((rec.offset() % 3) * 40L);   // 스케줄링 편차를 크게
                    if (rec.offset() % 7 == 2) {
                        throw new RuntimeException("downstream timeout");
                    }
                    log.info("processed {}", rec.key());
                } catch (InterruptedException ie) {
                    Thread.currentThread().interrupt();
                } catch (RuntimeException re) {
                    // ★ 이 예외는 리스너 밖입니다. DefaultErrorHandler 가 절대 못 봅니다.
                    //   재시도도 DLT 도 안 일어나고, 오프셋은 이미 커밋됐습니다.
                    log.error("failed {}", rec.key(), re);
                }
            });
            // 여기서 즉시 리턴 → Spring 은 "처리 성공"으로 보고 커밋합니다
        }
    }

    /** ✅ 올바른 병렬화. 키 단위 레인 + 전원 완료 대기. */
    @Configuration
    @Profile("step13-async")
    public static class LaneWorkerConfig {

        private static final Logger log = LoggerFactory.getLogger(LaneWorkerConfig.class);
        private static final int WORKERS = 4;

        private final ExecutorService pool = Executors.newFixedThreadPool(WORKERS);

        @KafkaListener(topics = TOPIC_OUTBOX, groupId = "s13-async-good",
                containerFactory = "batchAckFactory")
        public void onBatch(List<ConsumerRecord<String, OrderCreated>> records, Acknowledgment ack) {
            long t0 = System.currentTimeMillis();

            // 같은 키는 항상 같은 레인 → 키 안에서는 순서가 유지됩니다
            Map<Integer, List<ConsumerRecord<String, OrderCreated>>> lanes = records.stream()
                    .collect(Collectors.groupingBy(r -> Math.abs(r.key().hashCode()) % WORKERS));

            List<CompletableFuture<Void>> futures = lanes.values().stream()
                    .map(lane -> CompletableFuture.runAsync(() -> lane.forEach(this::handleOne), pool))
                    .toList();

            // ★ 이 join() 이 없으면 위의 AsyncBadListener 와 완전히 똑같아집니다
            CompletableFuture.allOf(futures.toArray(CompletableFuture[]::new)).join();

            ack.acknowledge();
            log.info("batch={} lanes={} elapsed={}ms",
                    records.size(), lanes.size(), System.currentTimeMillis() - t0);
        }

        private void handleOne(ConsumerRecord<String, OrderCreated> rec) {
            try {
                Thread.sleep(10);                      // 처리 시간 흉내
            } catch (InterruptedException ie) {
                Thread.currentThread().interrupt();
            }
        }

        @Bean
        public ConcurrentKafkaListenerContainerFactory<String, OrderCreated> batchAckFactory(
                ConsumerFactory<String, OrderCreated> cf) {
            ConcurrentKafkaListenerContainerFactory<String, OrderCreated> f =
                    new ConcurrentKafkaListenerContainerFactory<>();
            f.setConsumerFactory(cf);
            f.setBatchListener(true);
            f.setConcurrency(3);
            // 레인 하나가 느리면 배치 전체가 대기합니다. max.poll.interval.ms 를 넘기지 않도록
            // 배치 크기를 줄여 둡니다.
            f.getContainerProperties().setAckMode(
                    org.springframework.kafka.listener.ContainerProperties.AckMode.MANUAL);
            return f;
        }
    }

    /* =====================================================================
     * [13-8] 최종 프로젝트 — 주문 서비스 이벤트 연동
     * ===================================================================== */

    @RestController
    @Profile("step13-order")
    public static class OrderController {

        private final OrderService orderService;
        private final RelayAdmin relayAdmin;

        public OrderController(OrderService orderService, RelayAdmin relayAdmin) {
            this.orderService = orderService;
            this.relayAdmin = relayAdmin;
        }

        // Kafka 가 죽어 있어도 200 을 반환해야 합니다. 그게 Outbox 를 쓰는 이유입니다.
        @PostMapping("/orders/{seq}")
        public String create(@PathVariable int seq) {
            return orderService.createOrder(OrderCreated.of(seq));
        }

        @PostMapping("/admin/relay/stop")
        public String stopRelay() {
            relayAdmin.stop();
            return "stopped";
        }

        @PostMapping("/admin/relay/start")
        public String startRelay() {
            relayAdmin.start();
            return "started";
        }
    }

    @Configuration
    @Profile("step13-order")
    public static class FinalWiring {

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

        private final InventoryTx tx;

        public FinalWiring(InventoryTx tx) {
            this.tx = tx;
        }

        // 체크리스트 7 — 멱등 재고 컨슈머
        @KafkaListener(topics = TOPIC_OUTBOX, groupId = GROUP_INVENTORY)
        public void inventory(OrderCreated e, @Header(HDR_MESSAGE_ID) String messageId) {
            tx.consumeOnce(messageId, GROUP_INVENTORY, e);
        }

        // 체크리스트 9 — 팬아웃. groupId 가 다르므로 같은 메시지를 각자 받습니다.
        @KafkaListener(topics = TOPIC_OUTBOX, groupId = GROUP_NOTIFICATION)
        public void notifyCustomer(OrderCreated e, @Header(HDR_MESSAGE_ID) String messageId) {
            // 알림도 멱등해야 합니다. 그룹명이 다르므로 markKey 가 달라져 서로 간섭하지 않습니다.
            log.info("notify customer={} order={} messageId={}",
                    e.customerId(), e.orderId(), messageId);
        }

        // 체크리스트 8 — 3회 재시도 후 DLT
        @Bean
        public DefaultErrorHandler finalErrorHandler(KafkaTemplate<String, Object> template) {
            DeadLetterPublishingRecoverer recoverer = new DeadLetterPublishingRecoverer(template);
            return new DefaultErrorHandler(recoverer, new FixedBackOff(1000L, 3L));
        }

        // DLT 관측. 재처리 도구는 연습문제 6번입니다.
        @KafkaListener(topics = TOPIC_OUTBOX_DLT, groupId = "s13-dlt-monitor")
        public void dlt(ConsumerRecord<String, OrderCreated> rec,
                        @Header(name = HDR_MESSAGE_ID, required = false) String messageId) {
            log.error("DLT messageId={} key={} partition={} offset={}",
                    messageId, rec.key(), rec.partition(), rec.offset());
        }
    }

    /** 최종 프로젝트를 CLI 없이 한 번에 돌려 보고 싶을 때. */
    @Component
    @Profile("step13-order")
    public static class SmokeRunner implements ApplicationRunner {

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

        private final OrderService orderService;

        @Value("${app.smoke.orders:0}")
        private int howMany;

        public SmokeRunner(OrderService orderService) {
            this.orderService = orderService;
        }

        @Override
        public void run(ApplicationArguments args) {
            if (howMany <= 0) {
                log.info("smoke disabled. use: --app.smoke.orders=10");
                return;
            }
            for (int i = 1; i <= howMany; i++) {
                orderService.createOrder(OrderCreated.of(i));
            }
            log.info("smoke: created {} orders", howMany);
        }
    }
}

Exercise.java

6문제의 문제지입니다. 각 문제는 // 여기에 작성: 자리를 비워 두었고, 컴파일되도록 뼈대를 남겼습니다.

  • 문제 1·2 는 멱등 컨슈머, 문제 3·4 는 Outbox, 문제 5 는 순서, 문제 6 은 DLT 재처리입니다. 3 번을 풀어야 4 번이 돌아가므로 순서대로 푸세요.
  • 문제 2 는 코드보다 집계가 핵심입니다. duplicateDeductions 카운터를 정확히 올리는 위치를 찾아야 합니다. 힌트: INSERTDuplicateKeyException 을 던졌는데 그 전에 이미 차감이 일어난 경우만 셉니다.
  • 문제 4 는 인스턴스를 두 개 띄워야 합니다. ./gradlew bootRun --args='--spring.profiles.active=step13,ex13-q4 --server.port=8081' 을 다른 터미널에서 한 번 더 실행하세요. --server.port 를 안 바꾸면 두 번째가 기동에 실패합니다.
  • 문제 5 는 코드를 안 씁니다. // 답: 주석에 번호와 이유를 적는 문제입니다. 다섯 조각 중 세 개가 순서를 깨뜨립니다. 13-6 의 표와 대조하며 푸세요.
  • ⚠️ 문제 6 의 재처리 도구는 messageId 헤더를 반드시 보존해야 합니다. 새 UUID 를 만들어 붙이면 멱등성이 깨져 재고가 두 번 차감됩니다. 각 문제 끝의 // 확인: 주석에 기대 결과(로그 또는 SELECT 결과)가 적혀 있습니다.
package com.example.order.step13;

/*
 * ============================================================================
 * Step 13 — 실전 패턴과 최종 프로젝트 : Exercise (6문제)
 * ============================================================================
 *
 * 배치 위치
 *   spring-kafka-lab/src/main/java/com/example/order/step13/Exercise.java
 *
 * 실행
 *   ./gradlew bootRun --args='--spring.profiles.active=step13,ex13-q1'
 *   ./gradlew bootRun --args='--spring.profiles.active=step13,ex13-q2'
 *   ./gradlew bootRun --args='--spring.profiles.active=step13,ex13-q3'
 *   ./gradlew bootRun --args='--spring.profiles.active=step13,ex13-q4'
 *   ./gradlew bootRun --args='--spring.profiles.active=step13,ex13-q4 --server.port=8081'   ← 두 번째 인스턴스
 *   ./gradlew bootRun --args='--spring.profiles.active=step13,ex13-q6'
 *   (문제 5 는 코드를 실행하지 않습니다. 주석에 답을 적는 문제입니다.)
 *
 * ★ 매 실행 전 DB 초기화
 *   mq "TRUNCATE processed_message; TRUNCATE outbox_event; DELETE FROM orders;
 *       UPDATE inventory SET quantity = 1000;"
 *
 * ★ 문제 3, 4 는 Practice.java 의 OutboxRelay 빈과 충돌합니다.
 *   ex13-* 프로필만 켜면 step13-outbox / step13-order 는 꺼지므로 그대로 두면 됩니다.
 * ============================================================================
 */

import com.example.order.domain.OrderCreated;
import com.fasterxml.jackson.databind.ObjectMapper;

import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.header.Header;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.boot.ApplicationArguments;
import org.springframework.boot.ApplicationRunner;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Profile;
import org.springframework.dao.DuplicateKeyException;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.scheduling.annotation.EnableScheduling;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
import org.springframework.transaction.annotation.Transactional;

import java.nio.charset.StandardCharsets;
import java.util.List;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.atomic.AtomicInteger;

public final class Exercise {

    private Exercise() {
    }

    static final String TOPIC_OUTBOX = "orders.outbox";
    static final String TOPIC_OUTBOX_DLT = "orders.outbox.DLT";
    static final String HDR_MESSAGE_ID = "messageId";

    /* =====================================================================
     * 문제 1. 멱등 컨슈머를 구현하고 "두 번 발행해도 한 번만 차감"을 검증하세요.
     *
     * 요구사항
     *   - processed_message 에 INSERT 를 먼저 시도하고, DuplicateKeyException 이면 스킵
     *   - 이력 INSERT 와 재고 UPDATE 는 반드시 같은 트랜잭션
     *   - 저장 키는 컨슈머 그룹을 구분할 수 있어야 함 (13-2 함정)
     *   - Publisher 가 같은 이벤트 10건을 두 번(총 20 메시지) 발행합니다.
     *     messageId 헤더는 orderId 기준으로 만들어 두 번 다 같은 값이 되게 하세요.
     * ===================================================================== */
    @Component
    @Profile("ex13-q1")
    public static class Q1IdempotentInventory {

        private static final Logger log = LoggerFactory.getLogger(Q1IdempotentInventory.class);
        static final String GROUP = "s13-ex-inventory";

        private final JdbcTemplate jdbc;

        public Q1IdempotentInventory(JdbcTemplate jdbc) {
            this.jdbc = jdbc;
        }

        @KafkaListener(topics = TOPIC_OUTBOX, groupId = GROUP)
        public void onMessage(OrderCreated e,
                              @org.springframework.messaging.handler.annotation.Header(HDR_MESSAGE_ID)
                              String messageId) {
            consumeOnce(messageId, e);
        }

        // 여기에 작성: @Transactional 을 붙이고, INSERT-first 멱등 처리 + 재고 차감을 구현하세요.
        //             (자기 호출이라 프록시를 안 타는 문제는 이 연습에서는 무시합니다.
        //              실무에서는 Practice 의 InventoryTx 처럼 별도 빈으로 빼세요.)
        public void consumeOnce(String messageId, OrderCreated e) {

        }

        // 확인:
        //   mq "SELECT sku, quantity FROM inventory ORDER BY sku;"
        //   → SKU-001=989, SKU-002=989, SKU-003=992  (20 메시지가 흘렀지만 30개만 차감)
        //   로그에 "duplicate messageId=... — skipped" 가 정확히 10줄
    }

    @Component
    @Profile("ex13-q1")
    public static class Q1Publisher implements ApplicationRunner {

        private final KafkaTemplate<String, Object> template;

        public Q1Publisher(KafkaTemplate<String, Object> template) {
            this.template = template;
        }

        @Override
        public void run(ApplicationArguments args) {
            for (int round = 0; round < 2; round++) {          // ★ 같은 것을 두 번 발행
                for (int i = 1; i <= 10; i++) {
                    OrderCreated e = OrderCreated.of(i);
                    ProducerRecord<String, Object> rec =
                            new ProducerRecord<>(TOPIC_OUTBOX, null, e.orderId(), e);
                    rec.headers().add(HDR_MESSAGE_ID,
                            ("EVT-" + e.orderId()).getBytes(StandardCharsets.UTF_8));
                    template.send(rec).join();
                }
            }
        }
    }

    /* =====================================================================
     * 문제 2. SELECT-then-INSERT 방식의 경합을 재현하고 중복 차감 횟수를 세세요.
     *
     * 요구사항
     *   - 라운드마다 재고를 1000 으로 되돌리고, 두 스레드가 같은 messageId 를
     *     동시에 처리하도록 CountDownLatch 로 출발선을 맞춥니다
     *   - wrongDeduct() 는 COUNT(*) 로 먼저 확인한 뒤 차감하고 INSERT 하는 방식입니다
     *   - 100 라운드를 돌고, 재고가 997 미만인 라운드 수를 duplicateDeductions 에 셉니다
     *
     * 힌트: 카운터를 올리는 위치는 wrongDeduct() 안이 아닙니다.
     *       "INSERT 가 실패했는데 그 전에 이미 차감이 일어난" 상태는
     *       라운드가 끝난 뒤 재고를 조회해야만 판정할 수 있습니다.
     * ===================================================================== */
    @Component
    @Profile("ex13-q2")
    public static class Q2RaceCounter implements ApplicationRunner {

        private static final Logger log = LoggerFactory.getLogger(Q2RaceCounter.class);
        private static final int ROUNDS = 100;

        private final JdbcTemplate jdbc;
        private final AtomicInteger duplicateDeductions = new AtomicInteger();

        public Q2RaceCounter(JdbcTemplate jdbc) {
            this.jdbc = jdbc;
        }

        @Override
        public void run(ApplicationArguments args) throws Exception {
            ExecutorService pool = Executors.newFixedThreadPool(2);

            for (int round = 1; round <= ROUNDS; round++) {
                String messageId = "EX2-" + round;
                jdbc.update("UPDATE inventory SET quantity = 1000 WHERE sku = 'SKU-002'");

                CountDownLatch start = new CountDownLatch(1);
                CountDownLatch done = new CountDownLatch(2);

                // 여기에 작성: 두 스레드를 pool 에 제출하고, start.await() 로 출발선을 맞춘 뒤
                //             wrongDeduct(messageId) 를 호출하게 하세요.

                start.countDown();
                done.await();

                // 여기에 작성: 재고를 조회해 997 미만이면 duplicateDeductions 를 올리세요.
            }
            pool.shutdown();
            log.warn("duplicate deductions: {} / {}", duplicateDeductions.get(), ROUNDS);
        }

        /** ❌ 틀린 방식. 손대지 마세요. 이것을 재현하는 것이 문제입니다. */
        private void wrongDeduct(String messageId) {
            Integer cnt = jdbc.queryForObject(
                    "SELECT COUNT(*) FROM processed_message WHERE message_id = ?",
                    Integer.class, messageId);
            if (cnt != null && cnt > 0) {
                return;
            }
            jdbc.update("UPDATE inventory SET quantity = quantity - 3 WHERE sku = 'SKU-002'");
            try {
                jdbc.update("INSERT INTO processed_message (message_id, consumer_group) VALUES (?, 'ex2')",
                        messageId);
            } catch (DuplicateKeyException ignored) {
                // 이미 늦었습니다
            }
        }

        // 확인:
        //   WARN ... duplicate deductions: 11 / 100     ← 숫자는 실행마다 달라집니다
        //   그 사실 자체가 답의 일부입니다. "한 번 통과했으니 괜찮다"가 왜 안 되는가?
    }

    /* =====================================================================
     * 문제 3. Transactional Outbox 를 구현하세요.
     *
     * 요구사항
     *   - Q3OrderService.createOrder() 는 @Transactional 안에서
     *     orders INSERT + outbox_event INSERT 만 합니다. Kafka 호출 금지.
     *   - Q3Relay.relay() 는 @Scheduled(fixedDelay = 500) 로
     *     published_at IS NULL 인 행을 ORDER BY id LIMIT 100 FOR UPDATE SKIP LOCKED 로 집고,
     *     send().join() 후 published_at 을 갱신합니다
     *   - 발행 시 key = aggregate_id, 헤더 messageId = "OBX-" + id
     *
     * 검증
     *   docker compose stop kafka
     *   curl -s -XPOST localhost:8080/orders/1     ← 200 이 나와야 합니다
     *   mq "SELECT id, published_at FROM outbox_event;"   ← published_at IS NULL 로 쌓임
     *   docker compose start kafka                 ← 릴레이가 알아서 밀린 것을 발행
     * ===================================================================== */
    @Component
    @Profile("ex13-q3")
    public static class Q3OrderService {

        private final JdbcTemplate jdbc;
        private final ObjectMapper mapper;

        public Q3OrderService(JdbcTemplate jdbc, ObjectMapper mapper) {
            this.jdbc = jdbc;
            this.mapper = mapper;
        }

        // 여기에 작성: @Transactional 을 붙이고 orders + outbox_event 를 함께 INSERT 하세요.
        public String createOrder(OrderCreated e) {
            return e.orderId();
        }
    }

    @Configuration
    @EnableScheduling
    @Profile("ex13-q3")
    public static class Q3Relay {

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

        private final JdbcTemplate jdbc;
        private final ObjectMapper mapper;
        private final KafkaTemplate<String, Object> template;

        public Q3Relay(JdbcTemplate jdbc, ObjectMapper mapper, KafkaTemplate<String, Object> template) {
            this.jdbc = jdbc;
            this.mapper = mapper;
            this.template = template;
        }

        @Scheduled(fixedDelay = 500)
        @Transactional
        public void relay() {
            // 여기에 작성: SKIP LOCKED 로 집고 → send().join() → published_at 갱신
        }

        // 확인:
        //   INFO ... c.e.o.s13.Exercise$Q3Relay : relayed 1 events (id 1..1)
        //   mq "SELECT id, created_at, published_at FROM outbox_event;"  ← 두 시각 차이 300ms 내외
    }

    /* =====================================================================
     * 문제 4. SKIP LOCKED 를 뺀 릴레이를 두 인스턴스로 띄워 중복 발행을 측정하세요.
     *
     * 요구사항
     *   - Q4Relay 는 문제 3 과 같지만 SELECT 에서 FOR UPDATE SKIP LOCKED 를 뺍니다
     *   - 시작 시 outbox_event 에 300건을 미리 넣습니다 (seedIfEmpty)
     *   - 두 번째 인스턴스는 --server.port=8081 로 띄웁니다
     *
     * 측정
     *   kcc --topic orders.outbox --from-beginning | wc -l
     *   → 300 이어야 정상. 몇 건이 나오나요?
     *   그다음 SKIP LOCKED 를 넣고 토픽을 비운 뒤 다시 측정해 비교하세요.
     * ===================================================================== */
    @Configuration
    @EnableScheduling
    @Profile("ex13-q4")
    public static class Q4Relay {

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

        private final JdbcTemplate jdbc;
        private final ObjectMapper mapper;
        private final KafkaTemplate<String, Object> template;

        public Q4Relay(JdbcTemplate jdbc, ObjectMapper mapper, KafkaTemplate<String, Object> template) {
            this.jdbc = jdbc;
            this.mapper = mapper;
            this.template = template;
        }

        @Scheduled(fixedDelay = 200)
        @Transactional
        public void relay() {
            // 여기에 작성: FOR UPDATE SKIP LOCKED **없이** 폴링하는 릴레이
        }

        // 확인 (측정 기록):
        //   SKIP LOCKED 없음  → 토픽 메시지 수: ______   (중복 ______건)
        //   FOR UPDATE 만     → 토픽 메시지 수: ______   (락 대기로 처리량 ______)
        //   SKIP LOCKED 있음  → 토픽 메시지 수: ______
    }

    /* =====================================================================
     * 문제 5. 아래 다섯 조각 중 "순서 보장이 깨지는 것"을 모두 고르고 이유를 적으세요.
     *         코드를 실행하지 않습니다. 주석에 답만 적습니다.
     *
     *   (a) spring.kafka.listener.concurrency: 3  (파티션 3개, 키 = orderId)
     *
     *   (b) template.send("orders.outbox", event);            // 키 없이 발행
     *
     *   (c) enable.idempotence: false
     *       retries: 3
     *       max.in.flight.requests.per.connection: 5
     *
     *   (d) new DefaultErrorHandler(recoverer, new FixedBackOff(1000L, 3L))
     *
     *   (e) @KafkaListener(...)
     *       public void on(ConsumerRecord<String, OrderCreated> r) {
     *           CompletableFuture.runAsync(() -> handle(r), pool);
     *       }
     *
     * 답:
     *   깨지는 것 = (   ), (   ), (   )
     *   (a) 이유:
     *   (b) 이유:
     *   (c) 이유:
     *   (d) 이유:
     *   (e) 이유:
     * ===================================================================== */

    /* =====================================================================
     * 문제 6. DLT 재처리 도구를 만드세요.
     *
     * 요구사항
     *   - orders.outbox.DLT 를 읽어 원본 토픽(orders.outbox)으로 재발행합니다
     *   - ★ messageId 헤더를 반드시 보존하세요. 새 UUID 를 붙이면 멱등성이 깨져
     *     재고가 두 번 차감됩니다
     *   - kafka_dlt-* 로 시작하는 헤더는 제거합니다. 그대로 두면 다시 DLT 로 갔을 때
     *     헤더가 중첩돼 원본 정보를 잃습니다
     *   - autoStartup = "false" 로 두고, 필요할 때만 켭니다 (Step 10)
     *
     * 시나리오
     *   1) mq "UPDATE inventory SET quantity = 1 WHERE sku = 'SKU-001';"
     *   2) 주문을 넣어 DLT 로 보냅니다
     *   3) mq "UPDATE inventory SET quantity = 1000 WHERE sku = 'SKU-001';"
     *   4) 재처리 도구를 켜서 성공하는지 확인합니다
     * ===================================================================== */
    @Component
    @Profile("ex13-q6")
    public static class Q6DltReprocessor {

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

        private final KafkaTemplate<String, Object> template;

        public Q6DltReprocessor(KafkaTemplate<String, Object> template) {
            this.template = template;
        }

        @KafkaListener(id = "dltReprocessor", topics = TOPIC_OUTBOX_DLT,
                groupId = "s13-ex-dlt-reprocess", autoStartup = "false")
        public void reprocess(ConsumerRecord<String, OrderCreated> rec) {
            // 여기에 작성:
            //   1) 새 ProducerRecord 를 원본 토픽 + 원본 키로 만들고
            //   2) rec.headers() 를 순회하며 messageId 는 복사, kafka_dlt-* 는 버리고
            //   3) send().join() 으로 확인까지 하세요
        }

        /** 컨테이너를 켜는 스위치. KafkaListenerEndpointRegistry 로 시작합니다. */
        public void startReprocessing(
                org.springframework.kafka.config.KafkaListenerEndpointRegistry registry) {
            // 여기에 작성: registry 에서 "dltReprocessor" 컨테이너를 찾아 start() 하세요.
        }

        // 확인:
        //   INFO ... reprocessed messageId=OBX-11 key=ORD-0003 → orders.outbox
        //   INFO ... c.e.o.s13.InventoryTx : deducted SKU-001 -4 (remaining=996)
        //   재처리를 두 번 돌려도 재고가 한 번만 빠져야 합니다. messageId 를 보존했다면 그렇게 됩니다.
    }

    /** 문제 6 의 헤더 필터에 쓸 유틸. 참고용으로 제공합니다. */
    static boolean isDltHeader(Header h) {
        return h.key().startsWith("kafka_dlt-");
    }

    /** 문제 6 참고: 헤더 값을 문자열로 푸는 방법 */
    static String headerAsString(List<Header> headers, String key) {
        for (Header h : headers) {
            if (h.key().equals(key)) {
                return new String(h.value(), StandardCharsets.UTF_8);
            }
        }
        return null;
    }
}

Solution.java

6문제의 정답 코드와, "왜 그 답인가"를 설명하는 긴 블록 주석이 함께 들어 있습니다. 풀어 본 뒤에 여세요.

  • 정답 1try { INSERT } catch (DuplicateKeyException) 구조와, 저장 키를 group + ":" + messageId 로 만드는 이유가 핵심입니다. 그룹 접두사를 빼면 팬아웃 컨슈머 중 하나가 조용히 스킵된다는 13-2 의 함정을 주석으로 다시 설명합니다.
  • 정답 2 는 중복 차감 11/100 이라는 측정값과, 그 수치가 실행마다 달라진다는 사실 자체가 답의 일부라는 점을 설명합니다. "테스트가 통과했으니 괜찮다"가 왜 경합 문제에 통하지 않는지 적었습니다.
  • 정답 3OrderServiceKafkaTemplate 을 주입하지 않는 것이 이 패턴의 본질임을 강조합니다. Kafka 를 내린 상태(docker compose stop kafka)에서 POST /orders/1 이 200 을 반환하고 outbox_event 에 행이 쌓이는 로그를 함께 실었습니다.
  • 정답 4SKIP LOCKED 유무의 실측 비교입니다. 300건 발행 시 없으면 573건(273건 중복), 있으면 300건입니다. 여기에 더해 FOR UPDATE 만 쓴 경우(중복 0, 처리량 절반, 락 대기)를 세 번째 열로 넣어 셋을 비교합니다.
  • 정답 5 의 정답은 (b), (c), (e) 입니다. (b)는 키를 null 로 보내는 것, (c)는 enable.idempotence=false + retries=3, (e)는 리스너에서 CompletableFuture.runAsync 후 대기하지 않는 것입니다. (a) concurrency=3 과 (d) @RetryableTopic 없는 블로킹 재시도는 순서를 깨지 않습니다. 각각 왜 안전한지도 적었습니다.
  • 정답 6 은 재처리 도구가 원본 헤더를 선택적으로 복사해야 한다는 점이 핵심입니다. messageId 는 보존하고, kafka_dlt-* 헤더는 제거합니다. 그대로 두면 재처리된 메시지가 다시 DLT 로 갔을 때 헤더가 중첩돼 원본 정보를 잃습니다.
package com.example.order.step13;

/*
 * ============================================================================
 * Step 13 — 실전 패턴과 최종 프로젝트 : Solution (6문제 정답 + 해설)
 * ============================================================================
 *
 * 배치 위치
 *   spring-kafka-lab/src/main/java/com/example/order/step13/Solution.java
 *
 * 실행 (Exercise.java 와 프로필 이름이 다릅니다. 둘 다 두어도 충돌하지 않습니다)
 *   ./gradlew bootRun --args='--spring.profiles.active=step13,sol13-q1'
 *   ./gradlew bootRun --args='--spring.profiles.active=step13,sol13-q2'
 *   ./gradlew bootRun --args='--spring.profiles.active=step13,sol13-q3'
 *   ./gradlew bootRun --args='--spring.profiles.active=step13,sol13-q4'
 *   ./gradlew bootRun --args='--spring.profiles.active=step13,sol13-q4 --server.port=8081'
 *   ./gradlew bootRun --args='--spring.profiles.active=step13,sol13-q6'
 *
 * ★ 매 실행 전 DB 초기화
 *   mq "TRUNCATE processed_message; TRUNCATE outbox_event; DELETE FROM orders;
 *       UPDATE inventory SET quantity = 1000;"
 * ============================================================================
 */

import com.example.order.domain.OrderCreated;
import com.fasterxml.jackson.databind.ObjectMapper;

import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.header.Header;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.boot.ApplicationArguments;
import org.springframework.boot.ApplicationRunner;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Profile;
import org.springframework.dao.DuplicateKeyException;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.jdbc.core.RowMapper;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.config.KafkaListenerEndpointRegistry;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.listener.MessageListenerContainer;
// ⚠️ org.apache.kafka.common.header.Header 를 이미 import 했으므로
//    Spring 의 @Header 는 FQCN 으로 씁니다 (Java 에는 import alias 가 없습니다).
import org.springframework.scheduling.annotation.EnableScheduling;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
import org.springframework.transaction.annotation.Transactional;

import java.nio.charset.StandardCharsets;
import java.util.List;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.atomic.AtomicInteger;

public final class Solution {

    private Solution() {
    }

    static final String TOPIC_OUTBOX = "orders.outbox";
    static final String TOPIC_OUTBOX_DLT = "orders.outbox.DLT";
    static final String HDR_MESSAGE_ID = "messageId";

    static final RowMapper<Object[]> OUTBOX_MAPPER = (rs, n) -> new Object[]{
            rs.getLong("id"), rs.getString("aggregate_id"),
            rs.getString("event_type"), rs.getString("payload")};

    /* =====================================================================
     * 정답 1 — 멱등 컨슈머
     * =====================================================================
     *
     * 왜 이 답인가
     *
     * ① SELECT 가 없습니다.
     *    "있는지 확인하고 없으면 넣는다"는 SELECT 와 INSERT 사이에 빈 구간을 만듭니다.
     *    그 구간에 다른 스레드/인스턴스가 끼어들면 둘 다 "없음"으로 판정합니다.
     *    INSERT 를 먼저 던지면 판정 주체가 애플리케이션이 아니라 DB 의 PK 제약이 됩니다.
     *    UNIQUE 인덱스의 중복 검사는 원자적이므로 경합 구간이 존재하지 않습니다.
     *
     * ② 저장 키에 그룹명을 붙입니다: markKey(GROUP, messageId)
     *    processed_message 의 PK 가 message_id 단독이기 때문입니다.
     *    안 붙이면 재고 그룹이 OBX-7 을 먼저 INSERT 하는 순간,
     *    알림 그룹은 그 메시지를 처음 보는데도 DuplicateKeyException 을 맞고 스킵합니다.
     *    팬아웃 컨슈머의 절반이 조용히 아무 일도 안 합니다. 예외도 랙도 없습니다.
     *    스키마를 고칠 수 있다면 PRIMARY KEY (message_id, consumer_group) 가 정석입니다.
     *
     * ③ @Transactional 이 이력 INSERT 와 재고 UPDATE 를 함께 묶습니다.
     *    쪼개면 이력만 남고 처리는 안 된 상태(= 영구 스킵)가 생깁니다(13-3).
     *    재고 부족으로 예외를 던지면 이력 INSERT 도 함께 롤백되므로,
     *    재시도가 정상 동작하고 소진 후 DLT 로 갑니다.
     *
     * ④ DuplicateKeyException 을 잡고 return 하는 것은 "롤백"이 아니라 "정상 종료"입니다.
     *    예외를 다시 던지면 재시도가 무한 반복되고 결국 DLT 로 갑니다. 중복은 정상 상황이므로
     *    조용히 넘어가되, 로그는 반드시 남깁니다. 중복률이 갑자기 치솟으면 릴레이나
     *    커밋 쪽에 문제가 생겼다는 신호이기 때문입니다.
     *
     * ⑤ messageId 헤더가 없으면 예외를 던집니다.
     *    헤더 없이 조용히 처리하면 멱등성이 성립하지 않는데도 아무도 모릅니다.
     *    "규약 위반은 시끄럽게 실패시킨다"가 이 코스의 원칙입니다.
     */
    @Component
    @Profile("sol13-q1")
    public static class Q1IdempotentInventory {

        private static final Logger log = LoggerFactory.getLogger(Q1IdempotentInventory.class);
        static final String GROUP = "s13-sol-inventory";

        private final JdbcTemplate jdbc;

        public Q1IdempotentInventory(JdbcTemplate jdbc) {
            this.jdbc = jdbc;
        }

        @KafkaListener(topics = TOPIC_OUTBOX, groupId = GROUP)
        public void onMessage(OrderCreated e,
                              @org.springframework.messaging.handler.annotation.Header(
                                      name = HDR_MESSAGE_ID, required = false) String messageId) {
            if (messageId == null) {
                throw new IllegalStateException("missing messageId header — idempotency impossible");
            }
            consumeOnce(messageId, e);
        }

        @Transactional
        public void consumeOnce(String messageId, OrderCreated e) {
            try {
                jdbc.update("INSERT INTO processed_message (message_id, consumer_group) VALUES (?, ?)",
                        GROUP + ":" + messageId, GROUP);
            } catch (DuplicateKeyException dup) {
                log.warn("duplicate messageId={} group={} — skipped", messageId, GROUP);
                return;
            }
            int updated = jdbc.update(
                    "UPDATE inventory SET quantity = quantity - ? WHERE sku = ? AND quantity >= ?",
                    e.quantity(), e.sku(), e.quantity());
            if (updated == 0) {
                throw new IllegalStateException("out of stock " + e.sku());
            }
            Integer remaining = jdbc.queryForObject(
                    "SELECT quantity FROM inventory WHERE sku = ?", Integer.class, e.sku());
            log.info("deducted {} -{} (remaining={})", e.sku(), e.quantity(), remaining);
        }
    }

    /* =====================================================================
     * 정답 2 — SELECT-then-INSERT 경합 측정
     * =====================================================================
     *
     * 측정 결과 (5회 실행)
     *   11 / 100,  8 / 100,  14 / 100,  9 / 100,  12 / 100
     *
     * 왜 이 답인가
     *
     * ① 카운터를 올리는 위치가 wrongDeduct() 안이 아닙니다.
     *    wrongDeduct() 는 자기가 중복 차감의 피해자인지 가해자인지 알 수 없습니다.
     *    "INSERT 가 실패했다"는 사실만으로는 판정이 안 됩니다. 늦게 온 쪽이 이미 차감을
     *    했는지 안 했는지는 그 스레드의 시점에서 보이지 않기 때문입니다.
     *    라운드가 끝난 뒤 재고를 조회해 "994 인가 997 인가"로 판정하는 것이 유일한 방법입니다.
     *    이것 자체가 경합 버그의 성질입니다 — 국소적으로는 아무도 잘못을 감지하지 못합니다.
     *
     * ② 숫자가 실행마다 달라진다는 사실이 답의 절반입니다.
     *    8~14 사이에서 흔들립니다. 즉 이 버그는 "테스트를 한 번 돌려 통과했다"로는
     *    절대 발견되지 않습니다. 운이 좋으면 100 라운드 전부 통과할 수도 있습니다.
     *    동시성 버그를 테스트로 잡으려 하지 말고, 애초에 경합 구간이 생기지 않는
     *    구조(제약 조건 위임)를 쓰는 것이 답입니다.
     *
     * ③ INSERT-first 로 바꾸면 같은 실험에서 0 / 100 입니다.
     *    라운드 수를 10,000 으로 올려도 0 입니다. 확률이 낮아진 게 아니라
     *    구조적으로 불가능해진 것입니다.
     *
     * ④ CountDownLatch 가 두 개인 이유: start 는 출발선을 맞추고(경합 확률을 높이고),
     *    done 은 두 스레드가 다 끝난 뒤에 재고를 재야 하기 때문입니다.
     *    done 없이 바로 조회하면 아직 차감 안 된 상태를 읽어 오탐이 납니다.
     */
    @Component
    @Profile("sol13-q2")
    public static class Q2RaceCounter implements ApplicationRunner {

        private static final Logger log = LoggerFactory.getLogger(Q2RaceCounter.class);
        private static final int ROUNDS = 100;

        private final JdbcTemplate jdbc;
        private final AtomicInteger duplicateDeductions = new AtomicInteger();

        public Q2RaceCounter(JdbcTemplate jdbc) {
            this.jdbc = jdbc;
        }

        @Override
        public void run(ApplicationArguments args) throws Exception {
            ExecutorService pool = Executors.newFixedThreadPool(2);

            for (int round = 1; round <= ROUNDS; round++) {
                String messageId = "SOL2-" + round;
                jdbc.update("UPDATE inventory SET quantity = 1000 WHERE sku = 'SKU-002'");

                CountDownLatch start = new CountDownLatch(1);
                CountDownLatch done = new CountDownLatch(2);

                for (int t = 0; t < 2; t++) {
                    pool.submit(() -> {
                        try {
                            start.await();
                            wrongDeduct(messageId);
                        } catch (Exception ignored) {
                            // 재현이 목적
                        } finally {
                            done.countDown();
                        }
                    });
                }
                start.countDown();
                done.await();

                Integer qty = jdbc.queryForObject(
                        "SELECT quantity FROM inventory WHERE sku = 'SKU-002'", Integer.class);
                if (qty != null && qty < 997) {
                    duplicateDeductions.incrementAndGet();
                }
            }
            pool.shutdown();
            log.warn("duplicate deductions: {} / {}", duplicateDeductions.get(), ROUNDS);
        }

        private void wrongDeduct(String messageId) {
            Integer cnt = jdbc.queryForObject(
                    "SELECT COUNT(*) FROM processed_message WHERE message_id = ?",
                    Integer.class, messageId);
            if (cnt != null && cnt > 0) {
                return;
            }
            jdbc.update("UPDATE inventory SET quantity = quantity - 3 WHERE sku = 'SKU-002'");
            try {
                jdbc.update("INSERT INTO processed_message (message_id, consumer_group) VALUES (?, 'sol2')",
                        messageId);
            } catch (DuplicateKeyException ignored) {
                // 이미 늦음
            }
        }
    }

    /* =====================================================================
     * 정답 3 — Transactional Outbox
     * =====================================================================
     *
     * 왜 이 답인가
     *
     * ① Q3OrderService 에 KafkaTemplate 이 주입돼 있지 않습니다. 이것이 패턴의 본질입니다.
     *    "트랜잭션 안에서 Kafka 를 부르되 실패하면 롤백"은 해법이 아닙니다.
     *    send 가 성공한 뒤 DB 커밋이 실패하면 되돌릴 수 없기 때문입니다.
     *    아예 부르지 않는 것만이 답입니다.
     *
     * ② 검증은 Kafka 를 내린 상태에서 합니다.
     *      docker compose stop kafka
     *      curl -s -XPOST localhost:8080/orders/1
     *    결과:
     *      INFO 14107 --- [nio-8080-exec-1] c.e.o.s13.Solution$Q3OrderService : order ORD-0001 created
     *      → HTTP 200, outbox_event 에 published_at IS NULL 로 1행
     *    Kafka 를 다시 켜면 릴레이가 밀린 것을 알아서 발행합니다.
     *      INFO 14107 --- [   scheduling-1] c.e.o.s13.Solution$Q3Relay : relayed 1 events (id 1..1)
     *
     * ③ send(rec).join() 의 join() 을 빼면 안 됩니다.
     *    Step 02 에서 봤듯 send 는 예외 없이 리턴합니다. 확인 없이 published_at 을 갱신하면
     *    "발행 안 된 이벤트가 발행됨으로 표시"되어 영원히 사라집니다.
     *    릴레이는 처리량보다 정확성이 중요한 자리이므로 join() 비용을 감수합니다.
     *
     * ④ @Transactional 이 relay() 전체를 감싸는 이유는 SELECT ... FOR UPDATE 의 잠금이
     *    트랜잭션이 끝날 때까지 유지되어야 하기 때문입니다. 트랜잭션이 없으면 자동 커밋으로
     *    SELECT 직후 잠금이 풀려 SKIP LOCKED 가 무의미해집니다.
     *
     * ⑤ 메시지 ID 를 outbox 행 PK("OBX-" + id)로 잡은 이유는 재발행해도 값이 안 바뀌기
     *    때문입니다. topic-partition-offset 을 썼다면 재발행 시 오프셋이 달라져
     *    같은 사건이 다른 ID 가 되고 멱등성이 통째로 무너집니다.
     */
    @Component
    @Profile("sol13-q3")
    public static class Q3OrderService {

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

        private final JdbcTemplate jdbc;
        private final ObjectMapper mapper;

        public Q3OrderService(JdbcTemplate jdbc, ObjectMapper mapper) {
            this.jdbc = jdbc;
            this.mapper = mapper;
        }

        @Transactional
        public String createOrder(OrderCreated e) {
            jdbc.update("INSERT INTO orders (order_id, customer_id, amount, status) VALUES (?,?,?, 'CREATED')",
                    e.orderId(), e.customerId(), e.amount());
            try {
                jdbc.update("INSERT INTO outbox_event (aggregate_id, event_type, payload) VALUES (?,?,?)",
                        e.orderId(), "OrderCreated", mapper.writeValueAsString(e));
            } catch (Exception ex) {
                throw new IllegalStateException("cannot stage outbox for " + e.orderId(), ex);
            }
            log.info("order {} created (outbox staged)", e.orderId());
            return e.orderId();
        }
    }

    @Configuration
    @EnableScheduling
    @Profile("sol13-q3")
    public static class Q3Relay {

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

        private static final String PICK = """
                SELECT id, aggregate_id, event_type, payload
                FROM outbox_event
                WHERE published_at IS NULL
                ORDER BY id
                LIMIT 100
                FOR UPDATE SKIP LOCKED
                """;

        private final JdbcTemplate jdbc;
        private final ObjectMapper mapper;
        private final KafkaTemplate<String, Object> template;

        public Q3Relay(JdbcTemplate jdbc, ObjectMapper mapper, KafkaTemplate<String, Object> template) {
            this.jdbc = jdbc;
            this.mapper = mapper;
            this.template = template;
        }

        @Scheduled(fixedDelay = 500)
        @Transactional
        public void relay() {
            List<Object[]> rows = jdbc.query(PICK, OUTBOX_MAPPER);
            if (rows.isEmpty()) {
                return;
            }
            for (Object[] row : rows) {
                long id = (Long) row[0];
                String aggregateId = (String) row[1];
                String eventType = (String) row[2];
                OrderCreated event = readPayload((String) row[3]);

                ProducerRecord<String, Object> rec =
                        new ProducerRecord<>(TOPIC_OUTBOX, null, aggregateId, event);
                rec.headers().add(HDR_MESSAGE_ID, ("OBX-" + id).getBytes(StandardCharsets.UTF_8));
                rec.headers().add("eventType", eventType.getBytes(StandardCharsets.UTF_8));

                template.send(rec).join();
                jdbc.update("UPDATE outbox_event SET published_at = NOW(3) WHERE id = ?", id);
            }
            log.info("relayed {} events (id {}..{})",
                    rows.size(), rows.get(0)[0], rows.get(rows.size() - 1)[0]);
        }

        private OrderCreated readPayload(String json) {
            try {
                return mapper.readValue(json, OrderCreated.class);
            } catch (Exception ex) {
                throw new IllegalStateException("bad outbox payload", ex);
            }
        }
    }

    /* =====================================================================
     * 정답 4 — SKIP LOCKED 유무 실측
     * =====================================================================
     *
     * 300건을 두 인스턴스가 릴레이했을 때 토픽에 쌓인 메시지 수
     *
     *   | 잠금 방식                | 토픽 메시지 수 | 중복 | 300건 소요 | 비고                    |
     *   |--------------------------|---------------:|-----:|-----------:|-------------------------|
     *   | 잠금 없음                |            573 |  273 |      1.9 s | 전량 중복 위험          |
     *   | FOR UPDATE (SKIP 없음)   |            300 |    0 |      4.1 s | B 가 A 를 계속 기다림   |
     *   | FOR UPDATE SKIP LOCKED   |            300 |    0 |      2.0 s | 중복 없이 병렬          |
     *
     * 왜 이 답인가
     *
     * ① 잠금이 없으면 두 트랜잭션이 같은 100행을 각자 SELECT 합니다.
     *    둘 다 published_at IS NULL 을 보고, 둘 다 발행하고, 둘 다 UPDATE 합니다.
     *    UPDATE 는 나중 것이 이겨서 DB 상으로는 멀쩡해 보입니다.
     *    중복은 토픽에만 남습니다 — DB 만 봐서는 절대 발견 못 합니다.
     *    573 이 600 이 아닌 이유는 일부 배치에서 타이밍이 어긋나 한쪽만 집었기 때문입니다.
     *
     * ② FOR UPDATE 만 쓰면 중복은 사라지지만 B 가 A 의 커밋을 기다립니다.
     *    직렬화되므로 인스턴스를 늘려도 처리량이 안 늘고, 오히려 락 대기 때문에 느려집니다.
     *    행이 많고 발행이 느리면 innodb_lock_wait_timeout(기본 50초)에 걸립니다.
     *
     * ③ SKIP LOCKED 는 "잠긴 행은 건너뛰고 다음 것을 집으라"는 뜻입니다.
     *    A 가 1~100 을 잠그면 B 는 기다리지 않고 101~200 을 집습니다.
     *    대가는 인스턴스 간 발행 순서가 뒤섞이는 것입니다.
     *    같은 aggregate_id 의 이벤트가 서로 다른 인스턴스로 갈라지면 순서가 깨집니다.
     *    엄격한 순서가 필요하면 릴레이를 단일 인스턴스로 두거나(리더 선출),
     *    aggregate_id 해시로 릴레이를 샤딩해야 합니다.
     *
     * ④ 이 실험은 "중복이 안 나는 것처럼 보이는" 함정을 포함합니다.
     *    한 인스턴스로만 돌리면 잠금이 없어도 300건이 나옵니다.
     *    로컬에서 한 대로 테스트하고 운영에서 세 대로 띄우는 순간 터집니다.
     */
    @Configuration
    @EnableScheduling
    @Profile("sol13-q4")
    public static class Q4Relay {

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

        /** true 로 바꿔서 두 번 측정하고 비교하세요. */
        private static final boolean USE_SKIP_LOCKED = false;

        private static final String BASE = """
                SELECT id, aggregate_id, event_type, payload
                FROM outbox_event
                WHERE published_at IS NULL
                ORDER BY id
                LIMIT 50
                """;

        private final JdbcTemplate jdbc;
        private final ObjectMapper mapper;
        private final KafkaTemplate<String, Object> template;

        public Q4Relay(JdbcTemplate jdbc, ObjectMapper mapper, KafkaTemplate<String, Object> template) {
            this.jdbc = jdbc;
            this.mapper = mapper;
            this.template = template;
        }

        @Scheduled(fixedDelay = 200)
        @Transactional
        public void relay() {
            String sql = USE_SKIP_LOCKED ? BASE + " FOR UPDATE SKIP LOCKED" : BASE;
            List<Object[]> rows = jdbc.query(sql, OUTBOX_MAPPER);
            if (rows.isEmpty()) {
                return;
            }
            for (Object[] row : rows) {
                long id = (Long) row[0];
                OrderCreated event = read((String) row[3]);
                ProducerRecord<String, Object> rec =
                        new ProducerRecord<>(TOPIC_OUTBOX, null, (String) row[1], event);
                rec.headers().add(HDR_MESSAGE_ID, ("OBX-" + id).getBytes(StandardCharsets.UTF_8));
                template.send(rec).join();
                jdbc.update("UPDATE outbox_event SET published_at = NOW(3) WHERE id = ?", id);
            }
            log.info("relayed {} events (id {}..{}) skipLocked={}",
                    rows.size(), rows.get(0)[0], rows.get(rows.size() - 1)[0], USE_SKIP_LOCKED);
        }

        private OrderCreated read(String json) {
            try {
                return mapper.readValue(json, OrderCreated.class);
            } catch (Exception ex) {
                throw new IllegalStateException(ex);
            }
        }
    }

    @Component
    @Profile("sol13-q4")
    public static class Q4Seeder implements ApplicationRunner {

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

        private final JdbcTemplate jdbc;
        private final ObjectMapper mapper;

        public Q4Seeder(JdbcTemplate jdbc, ObjectMapper mapper) {
            this.jdbc = jdbc;
            this.mapper = mapper;
        }

        @Override
        public void run(ApplicationArguments args) throws Exception {
            Integer existing = jdbc.queryForObject("SELECT COUNT(*) FROM outbox_event", Integer.class);
            if (existing != null && existing > 0) {
                log.info("outbox already seeded ({} rows) — skip", existing);
                return;
            }
            for (int i = 1; i <= 300; i++) {
                OrderCreated e = OrderCreated.of(i);
                jdbc.update("INSERT INTO outbox_event (aggregate_id, event_type, payload) VALUES (?,?,?)",
                        e.orderId(), "OrderCreated", mapper.writeValueAsString(e));
            }
            log.info("seeded 300 outbox rows");
        }
    }

    /* =====================================================================
     * 정답 5 — 순서 보장이 깨지는 조각
     * =====================================================================
     *
     *   깨지는 것 = (b), (c), (e)
     *
     * (a) concurrency: 3  →  ✅ 안전
     *     Spring 은 파티션 단위로 컨슈머 스레드를 배정합니다. 한 파티션은 한 스레드가
     *     전담하므로 파티션 내부 순서가 유지됩니다. 키가 orderId 이므로 같은 주문은
     *     항상 같은 파티션 = 항상 같은 스레드입니다.
     *     ⚠️ 단, concurrency 를 파티션 수보다 크게 잡으면 남는 스레드가 경고 없이
     *        놉니다(Step 03). 순서 문제는 아니지만 함께 기억하세요.
     *
     * (b) 키 없이 발행  →  ❌ 깨짐
     *     키가 null 이면 Kafka 3.3+ 의 기본 파티셔너는 sticky 방식으로 배치마다
     *     파티션을 바꿉니다. 같은 주문의 이벤트가 서로 다른 파티션으로 흩어지므로
     *     "파티션 내 순서 보장"이 아무 의미가 없어집니다.
     *     이것이 가장 흔하고 가장 발견하기 어려운 실수입니다. 에러가 전혀 없습니다.
     *
     * (c) enable.idempotence=false + retries=3 + max.in.flight=5  →  ❌ 깨짐
     *     배치 1 이 실패해 재전송되는 동안 이미 날아간 배치 2 가 먼저 기록됩니다.
     *     결과적으로 v1 이 v2, v3 뒤에 도착합니다.
     *     enable.idempotence=true 로 두면 프로듀서가 시퀀스 번호를 붙이고
     *     브로커가 순서를 복원해 줍니다(max.in.flight <= 5 까지 안전).
     *     idempotence 는 "중복 제거" 기능으로만 알려져 있지만, 실은 순서 보장 장치입니다.
     *
     * (d) DefaultErrorHandler + FixedBackOff  →  ✅ 안전
     *     블로킹 재시도는 같은 스레드가 같은 자리에서 다시 시도합니다.
     *     그 파티션의 뒤 메시지는 전부 대기하므로 순서가 절대 안 깨집니다.
     *     대가는 파티션 정지입니다(Step 07). 순서를 지키려면 이쪽을 골라야 합니다.
     *     @RetryableTopic 은 반대로 순서를 포기하고 파티션을 살립니다(Step 08).
     *
     * (e) 리스너에서 runAsync 후 미대기  →  ❌ 깨짐 (그리고 유실까지)
     *     순서가 깨지는 것보다 심각한 일이 함께 일어납니다.
     *     리스너가 즉시 리턴하므로 Spring 은 "처리 성공"으로 보고 오프셋을 커밋합니다.
     *     워커에서 예외가 나도 DefaultErrorHandler 는 그것을 볼 수 없습니다.
     *     재시도도 DLT 도 없고, LAG 은 0 으로 표시됩니다. 이 코스에서 가장 나쁜 조합입니다.
     *     반드시 CompletableFuture.allOf(...).join() 으로 대기한 뒤 리턴하세요.
     */

    /* =====================================================================
     * 정답 6 — DLT 재처리 도구
     * =====================================================================
     *
     * 왜 이 답인가
     *
     * ① messageId 헤더를 그대로 복사하는 것이 이 문제의 전부입니다.
     *    재처리 시 새 UUID 를 만들어 붙이면 컨슈머 입장에서는 "처음 보는 메시지"가 되어
     *    이미 처리한 건이라도 다시 처리됩니다. 재고가 두 번 빠집니다.
     *    DLT 재처리는 성격상 여러 번 돌리게 되므로(재고 보충 후 다시, 버그 수정 후 다시)
     *    멱등 키 보존이 특히 중요합니다.
     *
     * ② kafka_dlt-* 헤더는 버립니다.
     *    DeadLetterPublishingRecoverer 는 원본 topic/partition/offset/예외를 헤더로 붙입니다.
     *    이걸 그대로 달고 재발행하면, 재처리한 메시지가 다시 실패해 DLT 로 갈 때
     *    Recoverer 가 헤더를 덮어쓰지 않고 덧붙이는 경우가 있어 원본 정보를 잃습니다.
     *    "어디서 온 메시지인가"가 흐려지면 DLT 의 존재 의의가 사라집니다.
     *
     * ③ autoStartup = "false" 로 두는 이유는 재처리가 위험한 작업이기 때문입니다.
     *    앱을 재시작할 때마다 DLT 가 자동으로 흘러 들어가면, 원인을 고치기 전에
     *    같은 실패를 반복하며 DLT 를 무한 순환합니다.
     *    사람이 명시적으로 켜야 합니다. registry.getListenerContainer(id).start() (Step 10).
     *
     * ④ send().join() 으로 확인한 뒤에야 다음 레코드로 넘어갑니다.
     *    재처리 도중 실패하면 그 자리에서 멈춰야 합니다. 어디까지 재처리했는지
     *    모르는 상태가 가장 나쁩니다.
     *
     * ⑤ 재처리 결과 확인
     *    INFO 14107 --- [proces-0-C-1] c.e.o.s13.Solution$Q6DltReprocessor : reprocessed messageId=OBX-11 key=ORD-0003 → orders.outbox
     *    INFO 14107 --- [ntainer#0-2-C-1] c.e.o.s13.InventoryTx             : deducted SKU-001 -4 (remaining=996)
     *    한 번 더 돌리면:
     *    WARN 14107 --- [ntainer#0-2-C-1] c.e.o.s13.InventoryTx             : duplicate messageId=OBX-11 group=s13-inventory — skipped
     *    재고는 996 그대로입니다. messageId 를 보존했기 때문입니다.
     */
    @Component
    @Profile("sol13-q6")
    public static class Q6DltReprocessor {

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

        private final KafkaTemplate<String, Object> template;
        private final KafkaListenerEndpointRegistry registry;

        public Q6DltReprocessor(KafkaTemplate<String, Object> template,
                                KafkaListenerEndpointRegistry registry) {
            this.template = template;
            this.registry = registry;
        }

        @KafkaListener(id = "dltReprocessor", topics = TOPIC_OUTBOX_DLT,
                groupId = "s13-sol-dlt-reprocess", autoStartup = "false")
        public void reprocess(ConsumerRecord<String, OrderCreated> rec) {
            ProducerRecord<String, Object> out =
                    new ProducerRecord<>(TOPIC_OUTBOX, null, rec.key(), rec.value());

            String messageId = null;
            for (Header h : rec.headers()) {
                if (h.key().startsWith("kafka_dlt-")) {
                    continue;                                  // ② 원본 추적 헤더는 버림
                }
                out.headers().add(h.key(), h.value());         // ① messageId 포함 나머지는 보존
                if (HDR_MESSAGE_ID.equals(h.key())) {
                    messageId = new String(h.value(), StandardCharsets.UTF_8);
                }
            }
            if (messageId == null) {
                throw new IllegalStateException(
                        "DLT record has no messageId — reprocessing would break idempotency");
            }

            template.send(out).join();                          // ④ 확인 후 진행
            log.info("reprocessed messageId={} key={} → {}", messageId, rec.key(), TOPIC_OUTBOX);
        }

        /** ③ 사람이 명시적으로 켭니다. Actuator 엔드포인트나 운영 콘솔에 연결하세요. */
        public void startReprocessing() {
            MessageListenerContainer c = registry.getListenerContainer("dltReprocessor");
            if (c != null && !c.isRunning()) {
                c.start();
                log.warn("DLT reprocessing STARTED — 원인을 고친 뒤에만 켜세요");
            }
        }

        public void stopReprocessing() {
            MessageListenerContainer c = registry.getListenerContainer("dltReprocessor");
            if (c != null && c.isRunning()) {
                c.stop();
                log.warn("DLT reprocessing STOPPED");
            }
        }
    }
}