문제는 “이중 쓰기”다
주문 서비스가 주문을 DB에 저장하고 OrderCreated를 Kafka에 발행하면 두 시스템에 쓰기가 발생합니다. 일반적인 @Transactional은 RDB만 묶을 뿐 Kafka까지 같은 원자적 경계로 만들지 않습니다.
DB만 성공
주문은 생성됐지만 이벤트가 없어 재고·배송 서비스가 주문을 알지 못합니다.
Kafka만 성공
이벤트는 전달됐지만 주문이 롤백되어 존재하지 않는 주문이 후속 처리됩니다.
Outbox는 도메인 변경과 “발행해야 할 사실”을 같은 DB 트랜잭션에 저장합니다. Inbox는 소비한 메시지와 업무 변경을 같은 트랜잭션에 저장합니다.
전체 아키텍처
| 계층 | 책임 | 책임지지 않는 것 |
|---|---|---|
| RDB Outbox | 도메인 변경과 이벤트 생성의 원자성 | Kafka 중복 제거 |
| Kafka | 전달·보관·재처리·파티션 내 순서 | 외부 RDB와 XA 원자성 |
| RDB Inbox | Consumer별 최종 중복 방지 | 동시 수정 충돌 |
| 낙관적 락 | lost update 방지 | 동일 메시지 재전달 방지 |
| Redis | hot dedupe와 lease | 영구 정합성의 근거 |
데이터 모델
Outbox
create table outbox_message (
id uuid primary key,
aggregate_type varchar(100) not null,
aggregate_id varchar(100) not null,
aggregate_version bigint not null,
event_type varchar(150) not null,
schema_version integer not null,
payload jsonb not null,
status varchar(20) not null,
attempts integer not null default 0,
available_at timestamp not null,
locked_until timestamp null,
lock_owner varchar(100) null,
occurred_at timestamp not null,
published_at timestamp null
);
create index ix_outbox_pending
on outbox_message(status, available_at);Inbox
create table inbox_message (
consumer_name varchar(150) not null,
message_id uuid not null,
payload_hash varchar(128),
status varchar(20) not null,
received_at timestamp not null,
processed_at timestamp,
primary key (consumer_name, message_id)
);consumer_name을 PK에 포함해야 같은 이벤트를 재고·알림·분석 Consumer가 각각 한 번씩 처리할 수 있습니다. 보존 기간은 Kafka 및 수동 replay 가능 기간보다 길게 잡습니다.
Spring 구현
Producer: 도메인과 이벤트를 한 번에 저장
@Transactional
public UUID createOrder(CreateOrder command) {
Order order = orderRepository.save(Order.create(command));
outboxRepository.save(OutboxMessage.create(
UUID.randomUUID(), "Order",
order.getId().toString(), order.getVersion(),
"OrderCreated", 1, OrderCreatedPayload.from(order)
));
return order.getId();
}이 트랜잭션 안에서 kafkaTemplate.send()를 원자성의 근거로 삼지 않습니다. Relay는 별도 프로세스나 스케줄 작업으로 실행합니다.
Relay: 짧게 claim하고 전송은 밖에서
List<OutboxMessage> batch = claimService.claim(100);
for (OutboxMessage message : batch) {
try {
kafkaTemplate.send(message.topic(),
message.aggregateId(), message.payload()).get();
claimService.markPublished(message.id());
} catch (Exception e) {
claimService.scheduleRetry(message.id(), e);
}
}Kafka 전송 후 완료 표시 전에 죽으면 중복 발행됩니다. 이것은 at-least-once의 정상 장애 창이며 Inbox가 흡수합니다. 대규모 시스템에서는 polling 대신 Debezium이 WAL/binlog를 읽게 할 수 있습니다.
Consumer: Inbox와 업무 변경을 한 번에 저장
@KafkaListener(topics = "order-events", groupId = "inventory-v1")
@Transactional
public void consume(OrderCreated event) {
int inserted = inboxRepository.insertIfAbsent(
"inventory-v1", event.eventId());
if (inserted == 0) return;
Inventory item = inventoryRepository
.findByProductId(event.productId()).orElseThrow();
item.reserve(event.quantity());
inboxRepository.markProcessed("inventory-v1", event.eventId());
}PostgreSQL에서는 INSERT ... ON CONFLICT DO NOTHING의 영향 행 수로 중복을 판정할 수 있습니다. DB 커밋 뒤 offset ACK 전에 장애가 나면 Kafka가 재전달하지만 Inbox PK가 재실행을 막습니다.
낙관적 락은 중복 제거가 아니다
@Entity
public class Inventory {
@Id private Long id;
@Version private Long version;
private int availableQuantity;
}@Version은 UPDATE 조건에 기존 version을 넣어 다른 트랜잭션의 선행 수정을 발견합니다. 동일 이벤트 재전달은 구분하지 못하므로 Inbox를 대신할 수 없습니다.
시간상의 중복
같은 message_id가 나중에 다시 도착하는 상황을 차단합니다.
동시성 충돌
서로 다른 요청이 같은 Aggregate를 동시에 덮어쓰는 것을 막습니다.
낙관적 락 충돌 시 Inbox INSERT도 함께 롤백되어야 합니다. 예외를 성공으로 삼키지 말고 제한된 재시도로 넘겨 최신 상태를 다시 읽습니다.
Redis는 보조 계층으로
빠른 중복 조회
SET inbox:{consumer}:{eventId} 1 NX EX 86400. Miss면 RDB Inbox를 확인합니다.
분산 claim
SET key token NX PX ttl. 해제 시 token 일치를 확인합니다.
영구 전송에 부적합
연결이 끊긴 Consumer의 메시지를 복구하지 않는 at-most-once 모델입니다.
소규모 Kafka 대안
Consumer Group, PEL, XACK, XAUTOCLAIM을 제공하지만 DB 부수효과에는 Inbox가 필요합니다.
Redis 키는 TTL·eviction·failover로 사라질 수 있습니다. 정확성의 기준은 RDB Inbox PK이며 Redis 장애 시 느려질 뿐 결과는 같아야 합니다.
Kafka 설정과 보장 범위
spring.kafka.producer.acks=all
spring.kafka.producer.properties.enable.idempotence=true
spring.kafka.consumer.enable-auto-commit=false
spring.kafka.consumer.properties.isolation.level=read_committedaggregate_id를 key로 사용해 같은 Aggregate를 같은 파티션으로 보냅니다.event_id,event_type,schema_version,aggregate_version,traceparent를 envelope에 둡니다.- 재시도에는 backoff와 최대 횟수를 두고 이후 DLT로 격리합니다.
- 외부 결제·메일 API에는
event_id기반 idempotency key를 전달합니다.
Kafka EOS는 Kafka read-process-write 범위에서 강력합니다. 외부 RDB까지 하나의 XA 트랜잭션으로 묶지는 않으므로 DB 부수효과에는 Inbox가 필요합니다.
장애가 나면 어떻게 되는가
| 장애 지점 | 복구 결과 |
|---|---|
| Producer DB 커밋 전 | 주문과 Outbox 모두 롤백 |
| DB 커밋 후 Kafka 발행 전 | Outbox가 남아 Relay/CDC가 복구 |
| Kafka 발행 후 완료 표시 전 | 중복 발행 가능, Inbox가 제거 |
| Consumer DB 커밋 전 | Inbox와 도메인 변경 모두 롤백 |
| DB 커밋 후 offset ACK 전 | 재전달되지만 Inbox가 중복 실행 차단 |
| 동일 Aggregate 동시 변경 | @Version 충돌 후 재시도 |
| Redis 장애 | RDB가 원본이면 성능만 저하 |
| 반복 실패 | 제한된 재시도 후 DLT 격리 |
운영 체크리스트
밀림과 지연
pending 수, oldest age, 발행 지연과 attempt 분포를 봅니다.
중복과 충돌
Inbox duplicate, 낙관적 락 충돌률, lag와 처리 시간을 기록합니다.
DLT와 Replay
원인과 재처리자를 기록하고 동일 event_id로 재투입합니다.
정리 정책
Inbox는 최대 replay 기간보다 오래 보존합니다.
권장 조합
대규모: PostgreSQL Outbox + Debezium CDC + Kafka + PostgreSQL Inbox + JPA @Version + Redis 보조 캐시 + Retry/DLT.
초기·중소규모: Spring Scheduled Polling + DB SKIP LOCKED/lease + Kafka + RDB Inbox.
공식 문서와 레퍼런스
- AWS — Transactional Outbox
이중 쓰기, 순서, 중복과 멱등 Consumer. - Debezium — Outbox Event Router
CDC 기반 Outbox와 Kafka 라우팅. - Spring Kafka — Transactions
Kafka Producer 트랜잭션. - Spring Kafka — Exactly Once
Kafka EOS의 적용 범위. - Spring — KafkaTransactionManager
Kafka와 RDB가 공유 XA가 아니라는 설명. - Redis — Pub/Sub
at-most-once 전달. - Redis — Streams
Consumer Group, ACK, PEL과 claim. - Microservices.io — Idempotent Consumer
Inbox/processed message 구현. - NServiceBus — Outbox
Message ID 중복 제거와 보존.