이전 글에서는 Kafka Consumer가 중복 메시지와 순서 역전을 어떻게 처리하는지 다뤘다. 그런데 Consumer를 아무리 견고하게 만들어도 메시지가 애초에 발행되지 않으면 소용없다.
이번 글은 반대쪽, Producer의 신뢰성 문제를 다룬다. "도메인 변경과 Kafka 발행을 어떻게 동시에 보장하는가."
문제: DB 쓰기와 Kafka 발행은 원자적이지 않다
유저가 username을 변경했다고 하자. DB를 업데이트하고 Kafka에 이벤트를 발행해야 한다. 순서를 어떻게 잡든 문제가 생긴다.
DB 먼저, Kafka 나중
userRepository.updateUsername(userId, newUsername) // DB 성공
kafkaTemplate.send("user.username.updated", payload) // Kafka 실패 → 이벤트 유실
DB엔 변경됐지만 다른 서비스는 모른다. Consumer는 아무리 잘 만들어도 오지 않는 메시지를 처리할 수 없다.
Kafka 먼저, DB 나중
kafkaTemplate.send("user.username.updated", payload) // Kafka 성공
userRepository.updateUsername(userId, newUsername) // DB 실패 → DB는 이전 상태
이벤트는 갔지만 DB엔 반영되지 않았다. Consumer는 존재하지 않는 변경을 처리하게 된다.
두 시스템에 걸친 원자적 커밋은 불가능하다. 분산 트랜잭션(2PC)으로 해결할 수도 있지만, 성능 비용과 복잡도가 크고 Kafka는 2PC를 지원하지 않는다.
해결: Outbox 테이블
핵심 아이디어는 단순하다. Kafka에 직접 쓰는 대신, 같은 DB 트랜잭션 안에 이벤트를 먼저 저장한다.
CREATE TABLE outbox_events (
id BIGSERIAL NOT NULL,
aggregate_id VARCHAR(100) NOT NULL,
aggregate_type VARCHAR(100) NOT NULL,
event_type VARCHAR(100) NOT NULL,
payload JSONB NOT NULL,
claimed_at TIMESTAMPTZ,
sent_at TIMESTAMPTZ,
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
CONSTRAINT pk_outbox_events PRIMARY KEY (id)
);
-- 미처리 이벤트만 빠르게 조회하기 위한 부분 인덱스
CREATE INDEX idx_outbox_unclaimed
ON outbox_events (created_at ASC, id ASC)
WHERE claimed_at IS NULL AND sent_at IS NULL;
도메인 변경과 이벤트 저장을 같은 트랜잭션에 묶는다.
@Transactional
fun updateUsername(userId: Long, newUsername: String) {
val newVersion = jdbcTemplate.queryForObject(
"""UPDATE users
SET username = ?, version = version + 1
WHERE user_id = ?
RETURNING version""",
Long::class.java, newUsername, userId
)!!
outboxRepository.save(
OutboxEvent(
aggregateId = userId.toString(),
aggregateType = "USER",
eventType = "USERNAME_UPDATED",
payload = """{"userId":$userId,"newUsername":"$newUsername","version":$newVersion}"""
)
)
// kafkaTemplate.send() 호출 없음 — Kafka 발행은 릴레이가 담당
}
트랜잭션이 커밋되면 DB 변경과 outbox 저장이 함께 반영된다. 트랜잭션이 롤백되면 둘 다 사라진다. Kafka 발행 실패가 도메인 로직에 영향을 주지 않는다.
Relay — Outbox에서 Kafka로
Outbox에 쌓인 이벤트를 주기적으로 읽어 Kafka에 발행하는 릴레이가 필요하다.
@Scheduled(fixedDelay = 1000, initialDelay = 5000)
fun relay() {
val events = outboxRepository.findAndClaim(100)
if (events.isEmpty()) return
for (event in events) {
try {
val topic = resolveTopic(event.eventType)
kafkaTemplate.send(topic, event.aggregateId, event.payload).get(5, TimeUnit.SECONDS)
outboxRepository.markSent(event.id)
} catch (ex: Exception) {
log.warn("Failed to publish outbox event id={}", event.id, ex)
outboxRepository.unclaim(event.id)
return
}
}
}
findAndClaim이 핵심이다. 단순 SELECT로 가져오면 여러 인스턴스가 동시에 같은 이벤트를 집어 중복 발행된다. FOR UPDATE SKIP LOCKED로 이를 막는다.
UPDATE outbox_events
SET claimed_at = NOW()
WHERE id IN (
SELECT id FROM outbox_events
WHERE claimed_at IS NULL
AND sent_at IS NULL
ORDER BY created_at, id
LIMIT ?
FOR UPDATE SKIP LOCKED
)
RETURNING id, aggregate_id, aggregate_type, event_type, payload, claimed_at, created_at
발행 성공 시 sent_at을 기록하고, 실패 시 claimed_at을 NULL로 돌려 다음 사이클에 재시도한다.
fun markSent(id: Long) {
jdbcTemplate.update(
"UPDATE outbox_events SET sent_at = NOW() WHERE id = ?", id
)
}
fun unclaim(id: Long) {
jdbcTemplate.update(
"UPDATE outbox_events SET claimed_at = NULL WHERE id = ?", id
)
}
예외 처리
Stale claim 복구
릴레이 인스턴스가 이벤트를 claim한 후 발행 직전에 죽으면 claimed_at은 세팅됐지만 sent_at은 없는 상태로 남는다. 이 이벤트는 다음 사이클에도 claimed_at IS NULL 조건에 걸리지 않아 영원히 처리되지 않는다.
@Scheduled(fixedDelay = 30000)
fun resetStaleClaims() {
jdbcTemplate.update(
"""UPDATE outbox_events
SET claimed_at = NULL
WHERE claimed_at < NOW() - INTERVAL '30 seconds'
AND sent_at IS NULL"""
)
}
30초가 지나도 sent_at이 없는 이벤트는 claim을 해제해 재처리 대상으로 돌린다.
처리 완료 레코드 정리
발행된 이벤트를 쌓아두면 테이블이 무한히 커진다.
@Scheduled(fixedDelay = 3600000)
fun deleteProcessed() {
jdbcTemplate.update(
"""DELETE FROM outbox_events
WHERE sent_at < NOW() - INTERVAL '7 days'"""
)
}
트레이드오프
at-least-once 보장 → Consumer 멱등성 필요
릴레이는 실패 시 재시도하므로 같은 이벤트가 두 번 발행될 수 있다. Outbox 패턴이 at-least-once를 만들어내는 구조이기 때문에, Consumer 측에서 반드시 멱등성 처리가 필요하다. 이전 글에서 다룬 processed_kafka_events 테이블이 바로
그 역할이다.
폴링 지연
fixedDelay = 1s면 최대 1초의 발행 지연이 생긴다. 실시간성이 중요한 도메인이라면 주기를 줄이거나 도메인 이벤트 저장 직후 릴레이를 즉시 트리거하는 방식을 고려할 수 있다.
DB 부하
outbox 테이블에 지속적인 read/write가 발생한다. 부분 인덱스(WHERE claimed_at IS NULL AND sent_at IS NULL)로 미처리 이벤트만 빠르게 조회할 수 있게 해두는 게 중요하다. 처리 완료 레코드를 주기적으로 지우는 것도 필수다.
마무리
두 글을 정리하면 하나의 흐름이 된다.
도메인 변경 + outbox 저장 → 릴레이가 Kafka 발행 → Consumer가 멱등성 처리
└─ Outbox 패턴 └─ 이전 글
(유실 방지) (중복/역전 방지)
Outbox가 "메시지는 반드시 보낸다" 를 보장하고, 멱등 Consumer가 "받더라도 한 번만 처리한다" 를 보장한다. 두 기법이 함께 있어야 end-to-end 신뢰성이 완성된다.
Outbox 패턴이 필요한 기준은 하나다. 메시지 유실이 비즈니스 문제로 이어지는 도메인인가. 알림, 결제, 사용자 동기화처럼 누락이 곧 데이터 불일치로 직결되는 경우라면 선택이 아닌 필수다.