Modu Square

Transactional Outbox Pattern 구현

문제

데이터 저장과 Kafka 이벤트 발행은 서로 다른 시스템에서 실행되므로 단일 트랜잭션 처리 불가

  • DB 커밋 후 Kafka 발행 실패: 데이터는 저장됐지만 다른 서비스가 변경 사실을 받지 못함
  • Kafka 발행 후 DB 커밋 실패: 존재하지 않는 데이터의 이벤트가 먼저 전달됨

해결 전략

  • 비즈니스 데이터와 이벤트를 같은 MySQL 트랜잭션에 저장해 함께 커밋하거나 롤백
  • 커밋 후 Kafka로 발행하고 전송 성공을 확인한 뒤 Outbox 삭제
  • 발행하지 못한 이벤트는 Outbox에 남겨 polling relay가 다시 발행
  • 재발행으로 같은 이벤트가 두 번 도착할 수 있어 소비 측에서 멱등 처리

기술 선택 이유

  • Two-Phase Commit - 제외

    • 참여 시스템의 커밋과 롤백을 조정할 수 있지만 MySQL과 Kafka 통합 제약과 처리 지연, coordinator 운영 부담 발생
  • CDC 기반 Outbox Relay - 제외

    • polling은 없앨 수 있지만 Debezium과 Kafka Connect의 실행 상태와 전송 지연, 장애 복구를 별도로 관리해야 함
  • Transactional Outbox - 채택

    • MySQL 트랜잭션 하나로 원자성이 보장돼 coordinator나 별도 CDC 구성 요소 없이 동작
    • Kafka 장애를 서비스 쓰기 요청과 분리하고 미전송 이벤트는 DB에 보존
Domain Service생성, 수정, 삭제같은 트랜잭션MySQL비즈니스 데이터 변경commit 또는 rollbackOutbox 이벤트 저장미전송 이벤트 보존커밋 후 발행Message Relay즉시 발행, 재시도성공 확인Kafka

구현

@TransactionalEventListener(phase = TransactionPhase.BEFORE_COMMIT)
public void createOutbox(OutboxEvent event) {
    outboxRepository.save(event.getOutbox());
}

@Async("messageRelayPublishEventExecutor")
@TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT)
public void publishEvent(OutboxEvent event) {
    publishEvent(event.getOutbox());
}

private void publishEvent(Outbox outbox) {
    try {
        messageRelayKafkaTemplate.send(
                outbox.getEventType().getTopic(),
                outbox.getPayload()
        ).get(1, TimeUnit.SECONDS);
        outboxRepository.delete(outbox);
    } catch (Exception e) {
        // polling relay가 다시 발행할 수 있도록 Outbox 유지
    }
}
  • BEFORE_COMMIT: 비즈니스 데이터와 Outbox를 함께 저장하고 실패 시 모두 롤백
  • AFTER_COMMIT: Kafka 발행을 시작하고 전송 성공을 확인한 뒤 Outbox 삭제
  • 남은 이벤트는 polling relay가 다시 발행
  • 소비 측은 토픽과 파티션, 게시글 단위로 마지막 처리 offset을 Redis에 기록하고 그보다 낮거나 같은 이벤트는 건너뜀

검증

초당 1건으로 60초간 게시글 생성. 실행 중 Kafka 20초 중단. 장애 구간의 쓰기 요청과 Outbox 저장 여부 확인

확인 항목결과
게시글 생성 요청61건
Outbox 최대 적재22건
최종 Outbox0건

Kafka 장애 중 Outbox에 미전송 이벤트가 쌓였다가 복구 후 0건으로 줄어든 Grafana 화면

Kafka 중단 뒤 미전송 이벤트 증가. 재시작 22초 후 0건으로 정상화

복구 후 정합성

최종 Outbox

SELECT COUNT(*)
FROM article.outbox;

결과: 0건

원본 게시글 수

SELECT COUNT(*)
FROM article.article
WHERE board_id = 9001;

결과: 100061

원본 count

SELECT article_count
FROM article.board_article_count
WHERE board_id = 9001;

결과: 100061

Kafka 복구 후 최종 정합성 확인

한계

  • Kafka 전송을 1초까지 동기로 기다려 Kafka가 느려지면 발행 스레드가 그만큼 묶임