정상 처리 비용
모든 주문에 Outbox row와 인덱스를 갱신하면 주문 트랜잭션의 DB 쓰기가 증가합니다.
Kafka Outbox Fallback · Visual Companion
정상 처리에서는 orders만 저장하고 Kafka로 직접 발행합니다. Kafka 발행에 실패하면 event_publications에 재발행 대상을 저장하고, 이 저장까지 실패하면 주문을 기준으로 발행 정보를 재구성합니다.
기존 Transactional Outbox는 주문과 발행 정보를 같은 트랜잭션에 저장합니다. Kafka-first Fallback은 정상 경로의 DB 입력을 줄이지만, 타임아웃과 fallback 저장 실패를 별도로 복구해야 합니다.
모든 주문에 Outbox row와 인덱스를 갱신하면 주문 트랜잭션의 DB 쓰기가 증가합니다.
응답이 지연됐을 뿐 Kafka가 이벤트를 저장했을 수 있습니다. 실패로 단정할 수 없습니다.
Kafka 발행과 event_publications 입력이 모두 실패하면 주문만 남습니다.
relay 또는 reconciler가 이미 전달된 eventId를 다시 발행할 수 있습니다.
첫 예제에 PostgreSQL, Kafka와 별도의 내구성 시스템을 추가해 핵심 절충안이 흐려집니다.
가장 강한 저장 보장을 유지하지만 정상 주문의 DB 입력을 줄이려는 목표를 충족하지 못합니다.
현재 Exposed·PostgreSQL 패턴을 재사용하면서 정상 처리와 복구 비용을 분리할 수 있습니다.
타임아웃과 reconciliation에서 발생할 수 있는 중복 발행을 처리합니다.
방식을 전환하면 DB 입력, Kafka 대기, 복구 기준이 함께 바뀝니다. DB 입력 수는 벤치마크가 아니라 현재 구현의 구조적 차이입니다.
Architecture Diagram은 시간 순서가 아니라 책임과 데이터 접근 관계를 보여줍니다. 단계별 시간 순서는 아래 장애 복구 시뮬레이션에서 확인합니다.
OrderController → PlaceOrderUseCase주문 요청PlaceOrderUseCase → TransactionalOrderWriter → orders주문 트랜잭션PlaceOrderUseCase → OrderEventPublisher → Kafka커밋 후 직접 발행OrderEventPublisher → EventPublicationRepository발행 실패 시 저장EventPublicationRelay → event_publications → Kafkaclaim 후 재발행PublicationReconciler → orders → event_publications누락 정보 재구성단계를 진행하면 현재 실행 계층과 누적 상태가 함께 변경됩니다. 타임아웃의 Kafka 이벤트 수는 0이 아니라 확인 불가로 표시합니다.
Claim은 같은 시점에 두 작업자가 같은 row를 처리하는 상황을 방지합니다. Kafka 발행 확인 후 DB 상태 변경 전에 프로세스가 중단되면 중복 발행은 여전히 가능합니다.
상태를 선택하면 진입 조건과 다음 상태를 확인할 수 있습니다. CLAIMED는 EventPublicationStatus enum 값이 아닙니다.
Class Diagram은 시나리오와 직접 연결되는 public 메서드와 사용 관계만 표시합니다.
orders 입력, eventId upsert, claim, SQL anti-join이 서로 다른 메서드에서 실행됩니다.
TransactionalOrderWriter.saveOrder()는 orders만 변경합니다. Kafka와 event_publications는 커밋 이후에 처리합니다.
OrderPlacedEvent.from(order)는 order-placed:{orderId}:v1을 사용합니다. 고유 event_id upsert가 재시도와 재구성을 단일 row로 모읍니다.
claimNextBatch()가 상태, nextAttemptAt, claim 만료, 정렬, batch limit을 DB에서 처리합니다.
findOrdersWithoutPublicationsCreatedOnOrBefore()가 cutoff와 anti-join으로 발행 정보가 없는 주문을 찾습니다.
애플리케이션은 PostgreSQL과 Kafka를 사용합니다. Demo admin endpoint는 기본적으로 비활성화되어 있습니다.
| 의미 | 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 |
| 동시 작업자 중 하나만 claim | concurrent 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 |
정상 경로에서는 orders만 입력하고 발행 정보를 별도로 저장하지 않습니다.
NOT_PUBLISHED row를 claim해 Kafka로 다시 발행합니다.
grace period 이후 orders를 기준으로 발행 정보를 재구성합니다.
Kafka가 이벤트를 받았는지 확정할 수 없습니다.
claim과 고정 eventId만으로 exactly-once 전달을 보장하지 않습니다.
같은 eventId를 다시 받아도 업무 결과가 중복되지 않게 구현해야 합니다.