AI Engineering Field Notes
Distributed Systems · Implementation Guide

Transactional
Outbox & Inbox

Spring Framework, Kafka, Redis와 낙관적 락을 조합해 DB 변경과 이벤트 전달 사이의 틈을 닫는 방법. “정확히 한 번”이라는 표현보다 장애가 발생해도 같은 결과로 수렴하는 구조를 설계합니다.

Spring BootKafkaRedisJPA @VersionDebezium CDCAt-least-once
DESIGN THESIS

RDB는 원자성을, Kafka는 전달과 순서를, Inbox는 멱등성을, 낙관적 락은 동시성 충돌을 맡는다. Redis는 이 정합성 모델을 빠르게 만드는 보조 계층이다.

문제는 “이중 쓰기”다

주문 서비스가 주문을 DB에 저장하고 OrderCreated를 Kafka에 발행하면 두 시스템에 쓰기가 발생합니다. 일반적인 @Transactional은 RDB만 묶을 뿐 Kafka까지 같은 원자적 경계로 만들지 않습니다.

FAILURE A

DB만 성공

주문은 생성됐지만 이벤트가 없어 재고·배송 서비스가 주문을 알지 못합니다.

FAILURE B

Kafka만 성공

이벤트는 전달됐지만 주문이 롤백되어 존재하지 않는 주문이 후속 처리됩니다.

Outbox는 도메인 변경과 “발행해야 할 사실”을 같은 DB 트랜잭션에 저장합니다. Inbox는 소비한 메시지와 업무 변경을 같은 트랜잭션에 저장합니다.

전체 아키텍처

확대·축소, 테마 전환, SVG/PNG 내보내기를 지원합니다.전체 화면 ↗
1. 원자 저장orders와 outbox를 같은 트랜잭션으로 커밋
2. RelayPolling 또는 Debezium CDC로 조회
3. Kafkaaggregate_id key로 개체별 순서 유지
4. 원자 소비inbox와 도메인 변경을 함께 커밋
5. 부수효과외부 API에 idempotency key 전달
계층책임책임지지 않는 것
RDB Outbox도메인 변경과 이벤트 생성의 원자성Kafka 중복 제거
Kafka전달·보관·재처리·파티션 내 순서외부 RDB와 XA 원자성
RDB InboxConsumer별 최종 중복 방지동시 수정 충돌
낙관적 락lost update 방지동일 메시지 재전달 방지
Redishot dedupe와 lease영구 정합성의 근거

데이터 모델

Outbox

create table outbox_message (
  id uuid primary key,
  aggregate_type varchar(100) not null,
  aggregate_id varchar(100) not null,
  aggregate_version bigint not null,
  event_type varchar(150) not null,
  schema_version integer not null,
  payload jsonb not null,
  status varchar(20) not null,
  attempts integer not null default 0,
  available_at timestamp not null,
  locked_until timestamp null,
  lock_owner varchar(100) null,
  occurred_at timestamp not null,
  published_at timestamp null
);
create index ix_outbox_pending
  on outbox_message(status, available_at);

Inbox

create table inbox_message (
  consumer_name varchar(150) not null,
  message_id uuid not null,
  payload_hash varchar(128),
  status varchar(20) not null,
  received_at timestamp not null,
  processed_at timestamp,
  primary key (consumer_name, message_id)
);

consumer_name을 PK에 포함해야 같은 이벤트를 재고·알림·분석 Consumer가 각각 한 번씩 처리할 수 있습니다. 보존 기간은 Kafka 및 수동 replay 가능 기간보다 길게 잡습니다.

Spring 구현

Producer: 도메인과 이벤트를 한 번에 저장

@Transactional
public UUID createOrder(CreateOrder command) {
    Order order = orderRepository.save(Order.create(command));
    outboxRepository.save(OutboxMessage.create(
        UUID.randomUUID(), "Order",
        order.getId().toString(), order.getVersion(),
        "OrderCreated", 1, OrderCreatedPayload.from(order)
    ));
    return order.getId();
}

이 트랜잭션 안에서 kafkaTemplate.send()를 원자성의 근거로 삼지 않습니다. Relay는 별도 프로세스나 스케줄 작업으로 실행합니다.

Relay: 짧게 claim하고 전송은 밖에서

List<OutboxMessage> batch = claimService.claim(100);
for (OutboxMessage message : batch) {
    try {
        kafkaTemplate.send(message.topic(),
            message.aggregateId(), message.payload()).get();
        claimService.markPublished(message.id());
    } catch (Exception e) {
        claimService.scheduleRetry(message.id(), e);
    }
}

Kafka 전송 후 완료 표시 전에 죽으면 중복 발행됩니다. 이것은 at-least-once의 정상 장애 창이며 Inbox가 흡수합니다. 대규모 시스템에서는 polling 대신 Debezium이 WAL/binlog를 읽게 할 수 있습니다.

Consumer: Inbox와 업무 변경을 한 번에 저장

@KafkaListener(topics = "order-events", groupId = "inventory-v1")
@Transactional
public void consume(OrderCreated event) {
    int inserted = inboxRepository.insertIfAbsent(
        "inventory-v1", event.eventId());
    if (inserted == 0) return;

    Inventory item = inventoryRepository
        .findByProductId(event.productId()).orElseThrow();
    item.reserve(event.quantity());
    inboxRepository.markProcessed("inventory-v1", event.eventId());
}

PostgreSQL에서는 INSERT ... ON CONFLICT DO NOTHING의 영향 행 수로 중복을 판정할 수 있습니다. DB 커밋 뒤 offset ACK 전에 장애가 나면 Kafka가 재전달하지만 Inbox PK가 재실행을 막습니다.

낙관적 락은 중복 제거가 아니다

@Entity
public class Inventory {
    @Id private Long id;
    @Version private Long version;
    private int availableQuantity;
}

@Version은 UPDATE 조건에 기존 version을 넣어 다른 트랜잭션의 선행 수정을 발견합니다. 동일 이벤트 재전달은 구분하지 못하므로 Inbox를 대신할 수 없습니다.

INBOX

시간상의 중복

같은 message_id가 나중에 다시 도착하는 상황을 차단합니다.

@VERSION

동시성 충돌

서로 다른 요청이 같은 Aggregate를 동시에 덮어쓰는 것을 막습니다.

낙관적 락 충돌 시 Inbox INSERT도 함께 롤백되어야 합니다. 예외를 성공으로 삼키지 말고 제한된 재시도로 넘겨 최신 상태를 다시 읽습니다.

Redis는 보조 계층으로

HOT DEDUPE

빠른 중복 조회

SET inbox:{consumer}:{eventId} 1 NX EX 86400. Miss면 RDB Inbox를 확인합니다.

RELAY LEASE

분산 claim

SET key token NX PX ttl. 해제 시 token 일치를 확인합니다.

PUB/SUB

영구 전송에 부적합

연결이 끊긴 Consumer의 메시지를 복구하지 않는 at-most-once 모델입니다.

STREAMS

소규모 Kafka 대안

Consumer Group, PEL, XACK, XAUTOCLAIM을 제공하지만 DB 부수효과에는 Inbox가 필요합니다.

Redis 키는 TTL·eviction·failover로 사라질 수 있습니다. 정확성의 기준은 RDB Inbox PK이며 Redis 장애 시 느려질 뿐 결과는 같아야 합니다.

Kafka 설정과 보장 범위

spring.kafka.producer.acks=all
spring.kafka.producer.properties.enable.idempotence=true
spring.kafka.consumer.enable-auto-commit=false
spring.kafka.consumer.properties.isolation.level=read_committed
  • aggregate_id를 key로 사용해 같은 Aggregate를 같은 파티션으로 보냅니다.
  • event_id, event_type, schema_version, aggregate_version, traceparent를 envelope에 둡니다.
  • 재시도에는 backoff와 최대 횟수를 두고 이후 DLT로 격리합니다.
  • 외부 결제·메일 API에는 event_id 기반 idempotency key를 전달합니다.

Kafka EOS는 Kafka read-process-write 범위에서 강력합니다. 외부 RDB까지 하나의 XA 트랜잭션으로 묶지는 않으므로 DB 부수효과에는 Inbox가 필요합니다.

장애가 나면 어떻게 되는가

장애 지점복구 결과
Producer DB 커밋 전주문과 Outbox 모두 롤백
DB 커밋 후 Kafka 발행 전Outbox가 남아 Relay/CDC가 복구
Kafka 발행 후 완료 표시 전중복 발행 가능, Inbox가 제거
Consumer DB 커밋 전Inbox와 도메인 변경 모두 롤백
DB 커밋 후 offset ACK 전재전달되지만 Inbox가 중복 실행 차단
동일 Aggregate 동시 변경@Version 충돌 후 재시도
Redis 장애RDB가 원본이면 성능만 저하
반복 실패제한된 재시도 후 DLT 격리

운영 체크리스트

METRICS

밀림과 지연

pending 수, oldest age, 발행 지연과 attempt 분포를 봅니다.

CONSUMER

중복과 충돌

Inbox duplicate, 낙관적 락 충돌률, lag와 처리 시간을 기록합니다.

RECOVERY

DLT와 Replay

원인과 재처리자를 기록하고 동일 event_id로 재투입합니다.

RETENTION

정리 정책

Inbox는 최대 replay 기간보다 오래 보존합니다.

권장 조합

대규모: PostgreSQL Outbox + Debezium CDC + Kafka + PostgreSQL Inbox + JPA @Version + Redis 보조 캐시 + Retry/DLT.

초기·중소규모: Spring Scheduled Polling + DB SKIP LOCKED/lease + Kafka + RDB Inbox.

공식 문서와 레퍼런스

  1. AWS — Transactional Outbox
    이중 쓰기, 순서, 중복과 멱등 Consumer.
  2. Debezium — Outbox Event Router
    CDC 기반 Outbox와 Kafka 라우팅.
  3. Spring Kafka — Transactions
    Kafka Producer 트랜잭션.
  4. Spring Kafka — Exactly Once
    Kafka EOS의 적용 범위.
  5. Spring — KafkaTransactionManager
    Kafka와 RDB가 공유 XA가 아니라는 설명.
  6. Spring Data — Optimistic Locking
    @Version 충돌 처리.
  7. Redis — Pub/Sub
    at-most-once 전달.
  8. Redis — Streams
    Consumer Group, ACK, PEL과 claim.
  9. Redis — Distributed Locks
    SET NX PX와 token 해제.
  10. Microservices.io — Idempotent Consumer
    Inbox/processed message 구현.
  11. NServiceBus — Outbox
    Message ID 중복 제거와 보존.