Kafka Outbox Fallback · Visual Companion

주문 트랜잭션의 Outbox 쓰기를 줄이면서 Kafka 장애 시 이벤트 발행을 복구한다

정상 처리에서는 orders만 저장하고 Kafka로 직접 발행합니다. Kafka 발행에 실패하면 event_publications에 재발행 대상을 저장하고, 이 저장까지 실패하면 주문을 기준으로 발행 정보를 재구성합니다.

1건정상 주문의 DB 입력
3회Kafka 직접 발행 시도
6개장애·복구 시나리오

모든 주문에 Outbox row를 저장하지 않으려면 복구 책임이 추가된다

기존 Transactional Outbox는 주문과 발행 정보를 같은 트랜잭션에 저장합니다. Kafka-first Fallback은 정상 경로의 DB 입력을 줄이지만, 타임아웃과 fallback 저장 실패를 별도로 복구해야 합니다.

정상 처리 비용

모든 주문에 Outbox row와 인덱스를 갱신하면 주문 트랜잭션의 DB 쓰기가 증가합니다.

불확실한 타임아웃

응답이 지연됐을 뿐 Kafka가 이벤트를 저장했을 수 있습니다. 실패로 단정할 수 없습니다.

두 번째 저장 실패

Kafka 발행과 event_publications 입력이 모두 실패하면 주문만 남습니다.

중복 발행 가능성

relay 또는 reconciler가 이미 전달된 eventId를 다시 발행할 수 있습니다.

구상 단계: 정상 처리 비용과 장애 복구 위험을 분리해 판단한다

검토 후 제외

Redis Streams를 대체 저장소로 사용

첫 예제에 PostgreSQL, Kafka와 별도의 내구성 시스템을 추가해 핵심 절충안이 흐려집니다.

검토 후 제외

정상 경로의 Outbox row를 유지

가장 강한 저장 보장을 유지하지만 정상 주문의 DB 입력을 줄이려는 목표를 충족하지 못합니다.

채택

실패한 발행 정보만 PostgreSQL에 저장

현재 Exposed·PostgreSQL 패턴을 재사용하면서 정상 처리와 복구 비용을 분리할 수 있습니다.

채택

고정된 eventId와 consumer 중복 처리 방지

타임아웃과 reconciliation에서 발생할 수 있는 중복 발행을 처리합니다.

기존 Transactional Outbox와 Kafka-first Fallback은 정상 처리 경로가 다르다

방식을 전환하면 DB 입력, Kafka 대기, 복구 기준이 함께 바뀝니다. DB 입력 수는 벤치마크가 아니라 현재 구현의 구조적 차이입니다.

정상 주문의 DB 입력
API 처리 중 Kafka 응답 대기
발행 정보 저장
누락 정보 재구성
consumer 중복 처리 방지

주문 저장, 직접 발행, 재발행, 발행 정보 재구성을 서로 다른 구성 요소가 담당한다

Architecture Diagram은 시간 순서가 아니라 책임과 데이터 접근 관계를 보여줍니다. 단계별 시간 순서는 아래 장애 복구 시뮬레이션에서 확인합니다.

API 계층

OrderController주문 API

애플리케이션 계층

PlaceOrderUseCase저장 후 발행 조정

저장 계층

TransactionalOrderWriterorders 트랜잭션
PostgreSQL.orders주문 데이터
EventPublicationRepository재발행 정보와 claim
PostgreSQL.event_publications장애·복구 대상

메시징 계층

OrderEventPublisherKafka 직접 발행과 fallback
Kafka.order-eventsOrderPlaced 이벤트

복구·운영 계층

EventPublicationRelayclaim 후 Kafka 재발행
PublicationReconciler누락 발행 정보 재구성

장애 조건을 선택해 DB, Kafka, API 상태가 바뀌는 과정을 확인한다

단계를 진행하면 현재 실행 계층과 누적 상태가 함께 변경됩니다. 타임아웃의 Kafka 이벤트 수는 0이 아니라 확인 불가로 표시합니다.

relay-aCLAIMED
relay-bEMPTY
API 결과
orders
event_publications
Kafka 이벤트
직접 발행 시도
relay 재시도
발행 상태
claim 소유자
API 계층애플리케이션 계층저장 계층메시징 계층복구·운영 계층

발행 상태와 claim 처리 상태를 구분한다

상태를 선택하면 진입 조건과 다음 상태를 확인할 수 있습니다. CLAIMED는 EventPublicationStatus enum 값이 아닙니다.

의미

진입 또는 다음 처리

실제 클래스가 시뮬레이션의 각 단계를 구현한다

Class Diagram은 시나리오와 직접 연결되는 public 메서드와 사용 관계만 표시합니다.

API 계층

OrderController

  • placeOrder(request): ResponseEntity<OrderResponse>
애플리케이션 계층

PlaceOrderUseCase

  • placeOrder(request): OrderResponse
저장 계층

TransactionalOrderWriter

  • saveOrder(...): OrderRecord
  • getOrder(orderId): OrderRecord
메시징 계층

OrderEventPublisher

  • publishDirectOrFallback(event): OrderPublicationStatus
저장 계층

EventPublicationRepository

  • upsertNotPublished(...)
  • claimNextBatch(...)
  • markPublished(...)
  • markRelayFailure(...)
  • findOrdersWithoutPublicationsCreatedOnOrBefore(...)
복구·운영 계층

EventPublicationRelay

  • scheduledRelay()
  • relayOnce(): RelayResult
복구·운영 계층

PublicationReconciler

  • scheduledReconcile()
  • reconcileOnce(): ReconcileResult
애플리케이션 계층

PublicationQueryService

  • findAll(): List<PublicationResponse>
복구·운영 계층

OutboxMetrics

  • recordDirectPublish(result)
  • recordFallbackStored(result)
  • recordRelay(result)
  • recordReconciler(result)

Exposed는 주문 트랜잭션과 재발행 정보 처리를 분리한다

orders 입력, eventId upsert, claim, SQL anti-join이 서로 다른 메서드에서 실행됩니다.

주문 트랜잭션

TransactionalOrderWriter.saveOrder()는 orders만 변경합니다. Kafka와 event_publications는 커밋 이후에 처리합니다.

고정된 eventId

OrderPlacedEvent.from(order)는 order-placed:{orderId}:v1을 사용합니다. 고유 event_id upsert가 재시도와 재구성을 단일 row로 모읍니다.

SQL claim

claimNextBatch()가 상태, nextAttemptAt, claim 만료, 정렬, batch limit을 DB에서 처리합니다.

누락 정보 조회

findOrdersWithoutPublicationsCreatedOnOrBefore()가 cutoff와 anti-join으로 발행 정보가 없는 주문을 찾습니다.

관련 서비스를 준비하고 정상 경로와 복구 경로를 확인한다

애플리케이션은 PostgreSQL과 Kafka를 사용합니다. Demo admin endpoint는 기본적으로 비활성화되어 있습니다.

테스트가 정상 처리, 장애 복구, 동시 claim과 정보 노출 방지를 검증한다

의미KafkaOutboxFallbackFlowTest
주문 트랜잭션이 orders만 변경transactional writer stores only order row
직접 발행 성공과 fallback row 부재placeOrder stores only order row and returns PUBLISHED_DIRECT when direct Kafka publish succeeds
세 번 실패 후 NOT_PUBLISHED 저장direct publish retries three times then stores NOT_PUBLISHED fallback row
제한 시간 내 반환과 fallback 저장direct publish timeout stores NOT_PUBLISHED fallback row
FALLBACK_STORE_FAILED 노출fallback insert failure returns FALLBACK_STORE_FAILED and records safe metric and log
재발행 후 PUBLISHED 전환relay publishes fallback row and marks it PUBLISHED
세 번째 실패 후 DEAD_LETTER 전환relay failure increments retry and moves to DEAD_LETTER
동시 작업자 중 하나만 claimconcurrent relay calls cannot claim the same row twice
만료된 claim의 재처리stale relay claim becomes eligible after claim ttl
DB에서 정렬과 batch limit 적용claimNextBatch applies SQL eligibility ordering and limit
고정된 eventId로 누락 정보 재구성reconciler reconstructs deterministic fallback row and documents duplicate risk
SQL cutoff와 anti-join 사용reconciler uses SQL cutoff and anti join for missing publications
원본 페이로드와 민감한 오류 정보 제외publication endpoint never exposes raw payload or raw exception text
직접 발행, fallback, relay, reconciler 메트릭metrics record direct failure fallback relay and reconciler outcomes

정상 처리의 DB 입력은 줄지만 중복과 복구 위험은 사라지지 않는다

해결

정상 주문의 Outbox 입력

정상 경로에서는 orders만 입력하고 발행 정보를 별도로 저장하지 않습니다.

해결

Kafka 장애 후 재발행

NOT_PUBLISHED row를 claim해 Kafka로 다시 발행합니다.

해결

발행 정보 저장 실패

grace period 이후 orders를 기준으로 발행 정보를 재구성합니다.

주의

타임아웃 이후 발행 여부

Kafka가 이벤트를 받았는지 확정할 수 없습니다.

주의

정확히 한 번 전달

claim과 고정 eventId만으로 exactly-once 전달을 보장하지 않습니다.

주의

consumer 중복 처리

같은 eventId를 다시 받아도 업무 결과가 중복되지 않게 구현해야 합니다.