Kafka를 걷어내고 DB 작업 큐로: 규모에 맞는 복잡도 되찾기

작성 · 수정


크레딧 선차감 + 비동기 이미지 생성 시스템에서 Kafka와 outbox 릴레이를 들어내고, 이미 존재하던 jobs 테이블을 그대로 작업 큐로 쓰도록 바꾼 기록이다. 프로덕션 코드 217줄과 테스트 389줄, 그리고 브로커 하나가 사라졌다.

문제

이 시스템에서 지키려고 하는 것은 하나다. 크레딧은 항상 정확하게 차감되고 환불된다.

클라이언트가 생성을 요청하면 멱등키를 확인하고, 조직 잔액을 차감하고(HOLD), job을 HOLDING 상태로 만들고, 원장에 기록한다. 여기까지가 동기 트랜잭션이다. 실제 이미지 생성은 건당 3~7초가 걸리는 외부 호출이라 비동기로 돌리고, 성공하면 CONFIRM, 실패하면 재시도하다가 시도 횟수를 소진하면 환불한다.

이 비동기 전달을 Kafka가 맡고 있었다. DB 커밋과 브로커 발행은 원자적이지 않다. 이 고전적인 이중 쓰기 문제를 막으려고 그 앞에 outbox 패턴이 붙어 있었다.

[HoldService] 한 트랜잭션
  1. idempotency_keys INSERT
  2. organization.balance 조건부 차감 UPDATE
  3. job(HOLDING) / ledger(HOLD) / outbox INSERT
        ↓
[OutboxRelay] 1초 폴링
  send().get(10s)로 브로커 ack를 확인한 뒤에만 markSent
        ↓
Kafka: generation-jobs 토픽 (파티션 3, + .DLT)
        ↓
[GenerationWorker] @KafkaListener(concurrency = 3)
  heartbeat 등록 → stub 호출 → confirm / fail

설계 자체는 좋다고 생각했는데 문제는 이 시스템에 그만한 규모가 오지 않는다는 것이었다.

브로커가 값을 하는 규모는 따로 있다. 유입이 처리 능력을 한참 앞질러 대기 물량이 수만 건씩 쌓이고, 그걸 디스크에 안전하게 눕혀둔 채 컨슈머 인스턴스를 늘려 파티션을 나눠 가져야 할 때다. 같은 이벤트를 각자 다른 속도로 읽어가는 소비자가 여럿이거나, 장애가 지나간 뒤 며칠치 로그를 되감아 재처리해야 하는 경우도 마찬가지다. 이런 조건이 하나라도 실재하면 브로커를 운영하는 비용은 아깝지 않다.

하지만 이 서비스에 그런 트래픽이 올 가능성은 낮았다. 그에 비해 브로커 하나를 정합성 있게 다루려고 쓴 코드는 너무 많았다. 오차를 막기 위한 코드가 계층을 이루고 있었다.

없애기 전에 무엇을 지고 있었는지 세어봤다.

프로덕션 클래스 7개 (217줄)

클래스존재 이유
OutboxEntry, OutboxRepository이중 쓰기를 막기 위한 중계 테이블
OutboxWriterjob 정보를 JSON으로 직렬화해 outbox에 적재
OutboxRelay미발송 건을 폴링해 발행하고 ack 확인 후 markSent
GenerationJobMessage브로커를 건너기 위한 메시지 DTO
KafkaTopicConfig토픽 + DLT 파티션 수 관리
KafkaConsumerConfigDefaultErrorHandler + DeadLetterPublishingRecoverer

테스트 5개 파일 (389줄) — OutboxRelayTest, OutboxRelayUnitTest, OutboxWriterTest, GenerationWorkerTest, GenerationWorkerDltTest. 뒤의 두 개는 @EmbeddedKafka로 브로커를 띄웠다.

그 외 — docker-compose.yml의 KRaft 브로커 서비스 16줄, spring.kafka 설정 블록(운영 15줄 + 테스트 12줄), 의존성 2개, 그리고 DeadJobSchedulerTask 안의 reapStaleHolding — 발행 실패로 영구히 HOLDING에 정체된 job을 회수하는 로직이었다.

마지막 항목이 특히 신호였다. 회수 로직이 존재하는 이유가 “브로커에 못 보냈을 수도 있어서”였다. 문제를 만들고, 그 문제를 막는 코드를 쓰고, 그 코드를 검증하는 테스트를 쓰고 있었다.

분석

지우기 전에 세 가지를 확인했다.

1. Kafka가 실제로 무엇을 해주고 있었나

  • 작업 전달 (프로듀서 → 컨슈머)
  • 리스너 컨테이너의 동시 소비 3
  • DLT를 통한 poison message 격리
  • 내구성 있는 로그와 리플레이

2. 그중 우리가 실제로 쓰던 것은

동시 소비 3과 DLT뿐이었다. 파티션 3은 형식만 갖췄을 뿐 컨슈머 그룹 재분배도, 다중 인스턴스 스케일아웃도 쓰지 않았다. 이벤트를 구독하는 다른 서비스가 없으니 팬아웃도, 리플레이도 필요 없었다. 처리량 상한은 브로커가 아니라 건당 3~7초짜리 외부 생성 API가 정했다. 초당 수천 건을 흘려보낼 수 있는 파이프 끝에 초당 0.3건짜리 수도꼭지가 달려 있었던 셈이다.

3. 그러면 큐가 DB 안에 있으면 무엇이 좋아지나

이게 결정적이었다. outbox 패턴이 존재하는 이유는 큐가 DB 밖에 있기 때문이다. 큐 자체가 DB 테이블이면 이중 쓰기 문제는 해결되는 게 아니라 정의부터 사라진다. 크레딧 차감과 큐 투입이 같은 트랜잭션의 두 INSERT면, 둘 중 하나만 성공하는 상태가 존재할 수 없다.

결정적 관찰이 하나 더 있었다. jobs 테이블은 이미 큐였다. status = HOLDING이 “대기 중”이고, attemptNo가 fencing 토큰이다. outbox는 jobs에 이미 있는 정보를 JSON으로 복제해 브로커로 보냈고, 컨슈머는 그걸 받아 다시 jobs를 조건부 UPDATE로 선점했다.

기존 Kafka 컨슈머의 첫 줄을 보면 명확하다.

@KafkaListener(topics = "${app.kafka.topic}", concurrency = "${app.kafka.partitions}")
public void consume(String payload) {
    GenerationJobMessage message = objectMapper.readValue(payload, GenerationJobMessage.class);

    int updated = jobRepository.startProcessingIfAttemptMatches(
            message.jobId(), message.attemptNo(), Instant.now());
    if (updated == 0) {
        log.info("무효한 메시지 무시: jobId={}, attemptNo={}", message.jobId(), message.attemptNo());
        return;  // ← Kafka가 무엇을 전달했든, 소유권은 여기서 결정된다
    }
    ...
}

소유권의 진실 원천은 이미 DB의 조건부 UPDATE였다. Kafka가 at-least-once로 중복 전달을 하든 말든 정확성은 DB가 지켰다. Kafka는 그 앞에 붙은 알림 채널에 지나지 않았다.

알림 채널을 없애고 DB를 직접 폴링해도 정확성 보장은 한 줄도 달라지지 않는다. 그렇다면 남는 건 무엇을 포기하느냐다.

포기하는 것판단
발행 즉시 소비 → 최대 500ms 폴링 지연건당 3~7초 작업에서 무의미한 크기
500ms마다 도는 SELECT 부하(status, id) 인덱스로 커버 가능
다중 인스턴스 스케일아웃조건부 UPDATE로 정확성은 유지되나 SKIP LOCKED 없이는 헛조회 경합. 단일 인스턴스 전제라 수용
DLT 격리max-attempts 3 상한이 같은 역할. 소진하면 환불로 종결되므로 큐를 잠식하지 않는다

전부 수용 가능했다.

해결

발행 경로를 지운다

HoldService와 RetryService에서 outboxWriter.write(...) 한 줄씩을 제거했다. hold 트랜잭션은 이제 멱등키 · 잔액 차감 · job · ledger 네 개만 커밋한다. 큐 투입은 job의 status = HOLDING 그 자체다.

컨슈머를 폴링 워커로 바꾼다

@Scheduled(fixedDelayString = "${app.scheduling.worker-interval-millis:500}")
public void dispatchPendingJobs() {
    List<Job> jobs = jobRepository.findByStatusOrderByIdAsc(JobStatus.HOLDING, PageRequest.of(0, batchSize));
    for (Job job : jobs) {
        if (!claim(job)) {
            continue;   // 다른 워커가 선점했거나 무효 → 다음 작업으로
        }
        if (!dispatch(job)) {
            return;     // 실행 슬롯이 없다 → 이번 주기 중단
        }
    }
}

선점은 그대로다. 바뀐 게 없다는 게 요점이다.

@Modifying(flushAutomatically = true, clearAutomatically = true)
@Query("""
        UPDATE Job j
        SET j.status = PROCESSING, j.updatedAt = :now
        WHERE j.id = :jobId
          AND j.status = HOLDING
          AND j.attemptNo = :attemptNo
        """)
int startProcessingIfAttemptMatches(...);

여러 워커가 같은 목록을 읽어도 괜찮다. 실제 소유권은 이 UPDATE의 반환 행 수가 결정한다. check-then-act 대신 확인과 실행을 하나의 조건부 UPDATE로 원자화한다 — 프로젝트의 이 원칙이 그대로 작업 큐 역할까지 겸하게 됐다.

되찾아야 했던 것: 동시성

여기서 한 번 넘어졌다. 처음 옮겼을 때 생성 작업이 스케줄러 스레드에서 직접 실행됐다. 결과가 두 가지였다.

  1. 스케줄러 풀이 작아 GenerationWorker와 DeadJobSchedulerTask가 같은 스레드를 두고 경합했다. 건당 3~7초짜리 외부 호출이 백로그만큼 이어지는 동안 회수 · 재시도 · 환불 스캔이 통째로 굶었다. 부하가 올라갈수록 복구 경로가 먼저 멈추는 구조였다.
  2. Kafka 리스너가 주던 동시 실행 3을 잃어 처리량이 1건으로 떨어졌다.

전용 executor로 실행을 분리했다.

@Bean("generationWorkerExecutor")
public ThreadPoolTaskExecutor generationWorkerExecutor(WorkerProperties workerProperties) {
    ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
    executor.setCorePoolSize(workerProperties.concurrency());
    executor.setMaxPoolSize(workerProperties.concurrency());
    executor.setQueueCapacity(0);   // ← 의도적으로 0
    executor.setThreadNamePrefix("generation-worker-");
    executor.initialize();
    return executor;
}

queueCapacity(0)이 핵심이다. 대기열 역할은 DB의 HOLDING이 한다. executor 안에 큐를 두면 이미 PROCESSING으로 선점됐지만 heartbeat는 등록되지 않은 작업이 쌓인다. 회수 스캐너 눈에는 그게 “죽은 작업”이다. 그래서 실행 슬롯을 확보하지 못하면 선점을 HOLDING으로 되돌리고 이번 주기를 끝낸다. 큐는 한 곳에만 있어야 한다.

reapStaleHolding을 지운다

이전에 HOLDING이 오래 머무는 건 “발행에 실패했다”는 뜻이었다. 그래서 시간 초과로 회수했다. 이제 HOLDING은 정상적인 대기 상태다. 경과 시간만으로 실패 처리하면 안 된다. 워커가 꺼져 있는 동안 쌓인 작업은 재기동 후 다음 폴링에서 그대로 처리된다.

Redis는 남겼다. 다만 역할이 바뀌었다 — 작업 전달 수단이 아니라, 살아있는 처리 작업을 확인하는 heartbeat 저장소로만 쓴다.

새로 생긴 실패 모드 하나

옮기고 나서 드러난 구멍이 있었다. 외부 생성은 성공했는데 결과를 DB에 반영(confirm)하는 데 실패하면? 초기 구현은 job을 FAILED로 바꿨다. 그러면 재시도가 이미 성공한 유료 외부 생성을 다시 실행했다.

FAILED 전이를 걷어내는 것만으로는 부족했다. heartbeat가 끊기면 timeout 회수 경로가 결국 재시도시키므로 중복 실행이 지연될 뿐이다. 이미 만들어진 resultUrl은 그대로 유실된다. 그래서 같은 호출 안에서 confirm을 3회 짧게 재시도해 순간적인 DB 장애에서 결과를 실제로 살린다. completeIfAttemptMatches가 status = PROCESSING과 attemptNo 일치를 조건으로 걸기 때문에 재시도가 중복 반영되지 않는다.

재시도를 소진하면 PROCESSING을 유지하고 회수 경로에 맡긴다. 이건 정직하게 트레이드오프로 남겼다 — 회수 후 재시도되면 외부 API가 두 번 호출된다. 크레딧 정확성은 유지되지만 외부 원가는 이중으로 든다. at-least-once를 택한 대가다.

인덱스

500ms마다 도는 폴링이 jobs를 풀스캔하고 있었다. 배치 폴링의 상태 필터와 ID 정렬을 함께 태우도록 (status, id) 인덱스를 추가했다.

결과

사라진 것

항목변화
Kafka/outbox 프로덕션 클래스7개 · 217줄 삭제
관련 테스트 코드389줄 삭제 (파일 5개)
전환 커밋 순증감+198 / −860
런타임 의존성spring-boot-starter-kafka 제거
로컬 실행 인프라MySQL + Redis + Kafka 브로커 → MySQL + Redis
데이터 모델outbox 테이블 제거
설정spring.kafka 블록 전체, app.kafka.*, app.holding.timeout-seconds 제거 / app.worker.batch-size, app.worker.concurrency 추가
테스트 인프라@EmbeddedKafka 전면 제거

사라진 실패 모드

코드 줄 수보다 이쪽이 본질이다.

  • outbox 발행 실패 → 없음
  • 브로커 ack 타임아웃 → 없음
  • 미발송으로 정체된 HOLDING 회수 → 없음 (HOLDING은 이제 정상 상태)
  • DLT 격리와 그 검증 → 없음 (max-attempts 상한이 종결시킨다)

남은 실패 모드는 두 개다. 워커 크래시(heartbeat 만료로 회수)와 confirm 실패(재시도 후 회수). 둘 다 하나의 회수 경로가 처리한다.

이중 쓰기를 막기 위한 코드가 아니라, 이중 쓰기가 필요 없는 구조가 됐다. 크레딧 차감 · job 등록 · 원장 기록이 하나의 RDB 트랜잭션이고, 작업 큐도 같은 DB 안에 있다.

그리고 정직하게, 한계

이 전환의 결론은 “Kafka가 틀렸다”가 아니다. 이 규모에서는 과했다다. 되돌아가야 할 신호는 분명하다.

  • 워커를 여러 인스턴스로 스케일아웃해야 할 때 — 정확성은 조건부 UPDATE가 지키지만 헛조회 경합이 커진다. 먼저 SKIP LOCKED, 그래도 안 되면 브로커
  • 다른 서비스가 같은 이벤트를 구독해야 할 때 — 팬아웃은 DB 큐가 못 하는 일이다
  • 이벤트 리플레이나 보존이 요구사항이 될 때
  • 폴링 부하가 DB 여유를 갉아먹기 시작할 때

지금 의도적으로 남긴 트레이드오프도 있다. 재시도 백오프가 없어 실패한 작업이 즉시 HOLDING으로 복귀하고, 단일 인스턴스를 전제하며, confirm 재시도를 소진하면 외부 API가 두 번 호출된다.

아키텍처 결정에 정답은 없고 유효기간만 있다. 지금 트래픽에 맞는 복잡도로 되돌리고, 되돌아갈 조건을 문서에 적어두는 것 — 그게 이번에 한 일이다.