아웃박스 패턴과 이벤트 유실 복구
카테고리 분류 처리 개선 중에 외부 API 장애로 유실되던 작업을 DB 아웃박스와 스케줄러로 복구하기
사용자가 처음 방문한 웹사이트는 외부 AI API를 통해 카테고리를 분류합니다.
기록 생성 API가 AI 응답을 기다리지 않도록, 초기 구현에서는 웹사이트를 저장한 뒤 Spring 이벤트를 발행하고 이벤트 리스너에서 분류를 수행했습니다.
@Transactional
fun upsertWebsite(command: UpsertWebsiteCommand): Website =
websiteRepository.save(Website(domain = command.domain))
.also {
eventPublisher.publishEvent(
WebsiteCategoryClassificationEvent(it.id!!)
)
}요청과 외부 API 호출을 분리하기에는 간단한 구조였습니다.
하지만 운영 데이터에서 생성된 지 오래된 웹사이트의 category가 계속 NULL로 남아 있는 현상을 발견했습니다.
메모리 이벤트가 사라지는 구간
Spring의 애플리케이션 이벤트는 메시지 브로커가 아닙니다.
이벤트가 처리되었다는 상태나 실패한 요청을 영속적으로 보관하지 않습니다.
다음 상황에서는 웹사이트만 저장되고 분류 요청은 사라질 수 있었습니다.
- 웹사이트 저장은 커밋되었지만 비동기 리스너가 실행되기 전에 인스턴스가 종료된다.
- 리스너가 외부 AI API를 호출하는 중 타임아웃이나 Rate Limit 오류가 발생한다.
- 재시도 중 애플리케이션이 재배포된다.
단순히 @Retryable을 붙이면 일시적인 장애에는 대응할 수 있습니다.
그러나 재시도 상태가 프로세스 메모리에만 있다면 인스턴스 종료 이후에는 복구할 수 없습니다.
반대로 category가 NULL인 웹사이트를 주기적으로 조회하는 방법은 구현은 쉽지만, 시도 횟수와 마지막 오류를 알 수 없고 의도적으로 미분류된 데이터와 실패한 작업을 구분하기 어렵습니다.
필요한 것은 비동기 실행 자체가 아니라 해야 할 작업을 잃지 않는 것이었습니다.
웹사이트와 작업을 함께 저장하기
웹사이트를 저장하는 트랜잭션에서 분류 요청을 나타내는 아웃박스도 함께 저장하도록 변경했습니다.
@Transactional
fun upsertWebsite(command: UpsertWebsiteCommand): Website =
with(command) {
val inserted =
websiteRepository.upsertByDomain(
Website(domain = domain, faviconUrl = faviconUrl)
)
websiteRepository
.getByDomain(domain)
.apply {
if (inserted) {
websiteCategoryClassificationOutboxRepository.save(
WebsiteCategoryClassificationOutbox(
websiteId = id!!
)
)
}
}
}CREATE TABLE website_category_classification_outbox (
id BINARY(16) PRIMARY KEY,
website_id BINARY(16) NOT NULL,
status ENUM(
'PENDING',
'PROCESSING',
'COMPLETED',
'FAILED'
) NOT NULL,
attempt_count INT NOT NULL,
attempted_at TIMESTAMP(6),
last_error_message VARCHAR(1000),
created_at TIMESTAMP(6) NOT NULL,
version BIGINT NOT NULL
);
CREATE UNIQUE INDEX uk_website_category_classification_outbox_website_id
ON website_category_classification_outbox (website_id);
CREATE INDEX idx_website_category_classification_outbox_polling
ON website_category_classification_outbox (status, attempted_at, created_at);웹사이트 저장과 아웃박스 저장 중 하나라도 실패하면 둘 다 롤백됩니다.
따라서 웹사이트는 존재하지만 분류 요청은 존재하지 않는 상태를 트랜잭션 경계에서 차단할 수 있습니다.
website_id에는 유니크 인덱스를 두었습니다.
동일한 웹사이트에 대해 여러 요청이 동시에 들어오더라도 분류 작업은 하나만 만들어집니다.
이벤트 리스너를 다시 붙이지 않은 이유
일반적인 아웃박스 구현에서는 커밋 직후 이벤트 리스너로 빠르게 처리하고, 스케줄러를 복구 수단으로 함께 사용하기도 합니다.
이번 기능에서는 이벤트 리스너를 제거하고 스케줄러만 두었습니다.
카테고리 분류는 호출량에 민감한 외부 AI API를 사용합니다.
새로운 웹사이트가 짧은 시간에 많이 등록될 때 이벤트를 즉시 처리하면 그 순간의 트래픽이 그대로 AI API 호출량이 됩니다.
반면 스케줄러가 일정한 간격으로 제한된 작업만 가져오면 DB의 아웃박스가 버퍼 역할을 합니다.
@Component
class WebsiteCategoryClassificationScheduler(
private val outboxService: WebsiteCategoryOutboxService
) {
@Scheduled(
fixedDelayString = "\${scheduler.website-category-classification.fixed-delay}"
)
fun processNextOutbox() {
outboxService.processNextOutbox()
}
}스케줄 한 번에 한 건만 처리하므로 인스턴스 하나를 기준으로 호출 속도를 예측할 수 있습니다.
다만 인스턴스가 여러 개라면 서로 다른 작업을 병렬로 처리할 수 있으므로, 전체 처리율은 인스턴스 수 × 실행 주기에 영향을 받습니다.
여러 인스턴스에서 하나의 작업 선점하기
두 인스턴스가 같은 PENDING 상태 행을 일반 SELECT로 읽은 뒤 각각 처리하면 AI API가 중복 호출됩니다.
작업을 가져오는 순간에는 짧은 트랜잭션 안에서 FOR UPDATE SKIP LOCKED를 사용했습니다.
override fun claimNext(
retryableAttemptedBefore: Instant,
recoverableAttemptedBefore: Instant
): WebsiteCategoryClassificationOutbox? =
dsl
.selectFrom(OUTBOX)
.where(
OUTBOX.STATUS.equal(PENDING)
.and(
OUTBOX.ATTEMPTED_AT.isNull()
.or(OUTBOX.ATTEMPTED_AT.lessOrEqual(retryableAttemptedBefore))
)
.or(
OUTBOX.STATUS.equal(PROCESSING)
.and(OUTBOX.ATTEMPTED_AT.lessOrEqual(recoverableAttemptedBefore))
)
)
.orderBy(OUTBOX.ATTEMPTED_AT.asc(), OUTBOX.CREATED_AT.asc())
.limit(1)
.forUpdate()
.skipLocked()
.fetchOneInto()FOR UPDATE는 선택한 행을 선점하고, SKIP LOCKED는 다른 트랜잭션이 이미 선점한 행을 기다리지 않고 다음 후보를 찾게 합니다.
그 결과 같은 아웃박스의 동시 선점은 막으면서 서로 다른 아웃박스는 여러 인스턴스가 병렬로 처리할 수 있습니다.
여기서 중요한 점은 DB 락을 AI 응답이 올 때까지 유지하지 않는 것입니다.
선점 트랜잭션에서는 상태를 PROCESSING으로 바꾸고 바로 커밋합니다.
val outbox =
transaction {
val now = Instant.now()
outboxRepository
.claimNext(
retryableAttemptedBefore = now - RETRY_DELAY,
recoverableAttemptedBefore = now - PROCESSING_TIMEOUT
)
?.let {
outboxRepository.save(
it.copy(
status = PROCESSING,
attemptCount = it.attemptCount + 1,
attemptedAt = now
)
)
}
} ?: return
websiteService.categorizeWebsite(
CategorizeWebsiteCommand(outbox.websiteId)
)외부 API 호출은 트랜잭션 밖에서 실행됩니다.
AI 응답을 기다리는 동안 행 락과 DB 커넥션을 점유하지 않기 위해서입니다.
실패 작업을 대기열 뒤로 보내기
처리에 실패한 작업은 attempted_at을 현재 시각으로 갱신하고 다시 PENDING 상태로 돌립니다.
아직 시도하지 않은 작업은 attempted_at이 NULL이므로 먼저 처리되고, 실패한 작업은 재시도 지연이 지난 뒤 다시 후보가 됩니다.
val status =
if (outbox.attemptCount >= MAX_ATTEMPT_COUNT) FAILED else PENDING
outboxRepository.save(
outbox.copy(
status = status,
attemptedAt = Instant.now(),
lastErrorMessage = exception.message
)
)이 방식으로 신규 작업을 먼저 처리하면서 실패한 작업을 논리적인 대기열 뒤로 보낼 수 있습니다.
최대 시도 횟수를 넘긴 작업은 FAILED 상태로 남기므로 운영자가 원인과 대상을 조회할 수 있습니다.
PROCESSING 상태 복구도 필요합니다.
인스턴스가 상태를 PROCESSING으로 커밋한 직후 종료되면 완료나 실패 상태로 바꿀 주체가 없어집니다.
그래서 attempted_at이 처리 제한 시간보다 오래된 PROCESSING 상태 작업은 다른 인스턴스가 다시 선점할 수 있게 했습니다.
exactly-once는 보장하지 않는다
SKIP LOCKED는 같은 순간의 중복 선점을 막지만 외부 API 호출의 exactly-once까지 보장하지는 않습니다.
예를 들어 AI API 호출은 성공했지만 결과를 저장하기 전에 인스턴스가 종료되면, 제한 시간이 지난 뒤 다른 인스턴스가 같은 작업을 다시 호출할 수 있습니다.
데이터베이스 트랜잭션과 외부 API를 하나의 원자적 연산으로 묶을 수 없기 때문입니다.
대신 최종 반영은 다음 조건으로 방어했습니다.
- 카테고리가 이미 존재하면 다시 분류하지 않는다.
COMPLETED상태의 아웃박스는 처리 후보에서 제외한다.- 작업 상태 변경에는 버전을 사용해 오래된 엔티티의 갱신을 감지한다.
외부 API가 멱등키를 지원한다면 호출 단계의 중복까지 더 강하게 방어할 수 있습니다.
현재 구조가 보장하는 것은 exactly-once 호출이 아니라 유실되지 않는 at-least-once 처리와 멱등한 최종 반영입니다.
결과
| 구분 | 실패 작업 복구 | 작업 영속화 | 실패 원인 추적 |
|---|---|---|---|
| 변경 전 | 불가능 | 없음 | 불가능 |
| 변경 후 | 가능 | 아웃박스 저장 | 가능 |
웹사이트 생성과 동시에 작업이 영속화되므로, 실패 원인과 시도 횟수도 데이터베이스에 남습니다.
이 작업에서 아웃박스는 단순한 이벤트 로그보다 복구 가능한 작업 큐에 가깝게 사용했습니다.
Kafka 같은 별도 메시지 시스템을 도입하지 않고도 현재 규모에서 필요한 내구성, 재시도, 다중 인스턴스 선점을 확보할 수 있었습니다.
구조는 메모리 이벤트보다 복잡해졌지만, 실패를 숨기는 대신 데이터로 남겨 다시 처리할 수 있게 된 점이 가장 큰 변화였습니다.