본문으로 건너뛰기
김신건의 로그

[Pattern] Outbox Pattern: DB + 메시지의 원자성

· 수정 · 📖 약 2분 · 582자/단어 #outbox #messaging #cdc #distributed #pattern
Outbox pattern, transactional outbox, Debezium, CDC, dual write problem, transactional messaging

정의

Outbox Pattern = DB 변경 + 메시지 발행원자성 보장. 이중 쓰기 (dual write) 문제 의 표준 해결.

문제: Dual Write

sequenceDiagram
    autonumber
    participant App
    participant DB
    participant Broker as Kafka/RabbitMQ

    App->>DB: BEGIN<br/>INSERT order<br/>COMMIT
    DB-->>App: OK
    App->>Broker: publish OrderCreated
    Note over App,Broker: ✗ App crash, 또는<br/>Broker timeout
    Broker--xApp: 응답 없음
    Note over DB,Broker: DB 는 INSERT 됨<br/>이벤트는 안 감 → *불일치*
시나리오결과
DB OK + Publish OK일관성
DB OK + Publish fail불일치 (DB 만 변경)
DB fail + Publish OK불일치 (이벤트만 갔음)
둘 다 failOK (취소)

두 시스템을 한 번에 일관 시키는 건 불가능. 2PC비현실적.

해결: Outbox 테이블

flowchart LR
    App -->|단일 트랜잭션| DB
    DB --> Orders[(orders)]
    DB --> Outbox[(outbox)]
    Outbox -.relay.-> Broker[(Broker)]
BEGIN;
  INSERT INTO orders (id, ...) VALUES (?, ...);
  INSERT INTO outbox (event_type, payload, aggregate_id)
    VALUES ('OrderCreated', ?, ?);
COMMIT;
-- 한 트랜잭션 = 원자성!

이후 별도 프로세스 가 outbox 를 읽고 broker 로 publish.

2가지 Relay 방식

1. Polling Publisher

flowchart LR
    Outbox[(outbox)] -.SELECT unpublished.-> Relay[Relay process]
    Relay -->|publish| Broker
    Relay -->|UPDATE published=true| Outbox
while True:
    rows = db.query("SELECT * FROM outbox WHERE published=false LIMIT 100")
    for row in rows:
        broker.publish(row.event)
        db.execute("UPDATE outbox SET published=true WHERE id=?", row.id)
    time.sleep(0.1)
  • 구현 간단.
  • DB 부담 (polling).
  • Latency ≈ polling 주기.

2. CDC (Change Data Capture)

flowchart LR
    DB[(DB)] -->|WAL / binlog| CDC[Debezium / Maxwell]
    CDC --> Kafka[(Kafka)]
    Kafka -->|topic| Consumers

DB 의 transaction log (PostgreSQL WAL / MySQL binlog)실시간 캡처.

도구DB
DebeziumPG, MySQL, MongoDB, Oracle 등
MaxwellMySQL
AWS DMS다양
Striim다양

IMPORTANT

Outbox + CDC2026 시점 마이크로서비스 메시징의 표준. 더 이상 애플리케이션 코드에서 publish 하지 않는다.

Outbox 테이블 디자인

CREATE TABLE outbox (
  id BIGSERIAL PRIMARY KEY,
  aggregate_type TEXT NOT NULL,        -- 'Order'
  aggregate_id TEXT NOT NULL,           -- 'order_42'
  event_type TEXT NOT NULL,             -- 'OrderCreated'
  payload JSONB NOT NULL,
  occurred_at TIMESTAMPTZ DEFAULT NOW(),
  published BOOLEAN DEFAULT false,
  published_at TIMESTAMPTZ
);

CREATE INDEX idx_outbox_unpublished
  ON outbox(id) WHERE published = false;

TIP

partial indexunpublished 만 빠른 SELECT.

처리 후 정리

flowchart TD
    P[Outbox 너무 큼?]
    P --> O1[옵션 1: published row 즉시 DELETE]
    P --> O2[옵션 2: 별도 archive 테이블]
    P --> O3["옵션 3: TTL (며칠 보관)"]

CDC 가 binlog stream 위주published 컬럼 자체 필요 없음. Debezium 이 INSERT 이벤트 그대로 전송.

메시지 순서 보장

flowchart LR
    Tx1[Tx 1: INSERT outbox id=100] --> Order
    Tx2["Tx 2: INSERT outbox id=99 (먼저 시작했지만 commit 늦음)"]
    Tx2 --> Order
    Note["DB ID 순 ≠ commit 순"]

CAUTION

concurrent transaction 때문에 id 순서 ≠ commit 순서. strict 순서 보장 이 필요하면 aggregate ID 별 partition 으로.

At-least-once + Idempotency

Outbox + relay 는 at-least-once. 중복 발행 가능 (relay 가 publish 후 update 전 crash).

consumer 가 idempotent 처리. 자세한 건 idempotency-keys.

Inbox Pattern (반대 방향)

flowchart LR
    Broker -->|consume| Inbox[(inbox)]
    Inbox -->|단일 트랜잭션| DB
    DB -->|이미 처리?| Skip[중복 skip]

수신 측의 멱등 보장. Outbox 와 대칭.

사용 상황

상황권장 방식
마이크로서비스 간 이벤트Outbox + CDC
레거시 DB (WAL 접근 불가)Polling Publisher
실시간 latency < 100msCDC 권장
단순 단일 서비스불필요 (직접 publish)
이벤트 감사 로그 필요Outbox + archive 테이블

Polling vs CDC 선택

flowchart TD
    S[Relay 방식 선택]
    S --> Q1{"DB WAL/binlog<br/>접근 가능?"}
    Q1 -->|"Yes"| Q2{"Latency 요구<br/>< 500ms?"}
    Q1 -->|"No"| Poll["Polling Publisher<br/>(구현 단순)"]
    Q2 -->|"Yes"| CDC["CDC (Debezium)<br/>(실시간, 저 latency)"]
    Q2 -->|"No"| Poll2["Polling Publisher<br/>(latency 허용)"]

TIP

Debezium + PostgreSQL WAL 조합이 2026 표준. wal_level=logical 설정 필요.

프레임워크 통합

Spring Boot + Spring Modulith

// @ApplicationModuleListener 가 Outbox 자동 관리
@ApplicationModuleListener
public void on(OrderCreated event) {
    // Spring Modulith 가 ApplicationEvent 를 Outbox 로 직렬화
}

Debezium Connector 설정

{
  "name": "outbox-connector",
  "config": {
    "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
    "database.hostname": "localhost",
    "table.include.list": "public.outbox",
    "transforms": "outbox",
    "transforms.outbox.type":
      "io.debezium.transforms.outbox.EventRouter",
    "transforms.outbox.table.field.event.id": "id",
    "transforms.outbox.table.field.event.key": "aggregate_id",
    "transforms.outbox.table.field.event.type": "event_type"
  }
}

운영 고려사항

1. Relay Lag 모니터링

-- 발행 안 된 건수 + 최고령 행 나이
SELECT
    COUNT(*) AS unpublished_count,
    EXTRACT(EPOCH FROM (NOW() - MIN(occurred_at))) AS oldest_sec
FROM outbox
WHERE published = false;

WARNING

oldest_sec 이 SLA (예: 30초) 를 넘으면 relay 가 막혔다는 신호. 알람 설정 필수.

2. Dead Letter 처리

ALTER TABLE outbox ADD COLUMN retry_count INT DEFAULT 0;
ALTER TABLE outbox ADD COLUMN last_error TEXT;

-- relay 에서 실패 시
UPDATE outbox
SET retry_count = retry_count + 1,
    last_error  = 'Broker timeout'
WHERE id = ?;

-- retry_count > 5 → DLQ 이동

3. 대용량 정리 정책

flowchart LR
    Pub[published 행 누적]
    Pub --> A{"보관 정책?"}
    A -->|"즉시"| Del["DELETE WHERE published=true<br/>(relay 직후)"]
    A -->|"TTL"| TTL["DELETE WHERE published_at &lt; NOW()-7d<br/>(배치 job)"]
    A -->|"Archive"| Arc["INSERT INTO outbox_archive<br/>+ DELETE FROM outbox"]
-- 효율적 배치 삭제 (lock 최소화)
DELETE FROM outbox
WHERE published = true
  AND published_at < NOW() - INTERVAL '7 days'
LIMIT 10000;

테스트 전략

# 통합 테스트: outbox 행 확인
def test_order_creates_outbox_event(db, client):
    resp = client.post('/orders', json={'item': 'Widget'})
    assert resp.status_code == 201

    rows = db.execute(
        "SELECT event_type, payload FROM outbox WHERE published=false"
    ).fetchall()
    assert len(rows) == 1
    assert rows[0]['event_type'] == 'OrderCreated'

# relay 테스트: 멱등성 확인
def test_relay_idempotent(db, broker, relay):
    db.execute("INSERT INTO outbox (event_type, payload) VALUES ('E', '{}')")
    relay.run_once()  # 첫 번째 relay
    relay.run_once()  # relay 가 crash 후 재시작 시나리오

    msgs = broker.consumed_messages('events')
    assert len(msgs) == 1  # 중복 없음 (consumer idempotent)

흔한 함정

WARNING

  1. *UPDATE published=truepublish 가 별도 transaction = relay crash 시 중복 발행. at-least-once 인정 + consumer idempotent.
  2. Outbox 의 순서 무시 = 다운스트림이 순서 가정. partition / sequence 명시.
  3. Outbox 가 너무 큼 = 정리 필요. 보관 정책.
  4. SELECT FOR UPDATE SKIP LOCKED 없이 multi worker = 같은 row 동시 처리 → 중복.

관련 위키

이 글의 용어 (6개)
[Distributed] Kafka: 분산 로그, partition, consumer groupdistributed-systems
정의 Apache Kafka = 분산 commit log. 고처리량 (수백만 msg/s), 영속, 수평 확장. event-driven 아키텍처 의 de facto. 핵심 개념: …
[Distributed] RabbitMQ: exchange, queue, routing keydistributed-systems
정의 RabbitMQ = AMQP 0.9.1 기반 traditional message broker. exchange → queue 라우팅, workload distribution…
[Pattern] Event Sourcing: 상태 대신 이벤트 누적distributed-systems
정의 Event Sourcing = 현재 상태 대신 과거 이벤트 시퀀스 를 저장. 현재 상태 = event 모두 replay 한 결과. 핵심 약속 - 과거 모든 변화 기록. 감사…
[Pattern] Idempotency Keys: 중복 요청 안전 처리distributed-systems
정의 Idempotency = 같은 요청을 N번 보내도 결과가 1번과 동일. 분산 시스템 / 결제 / API 의 안전망. [!IMPORTANT] 네트워크는 항상 timeout /…
[Pattern] Saga: 분산 트랜잭션의 보상 흐름distributed-systems
정의 Saga = 여러 service 의 로컬 트랜잭션 을 연결한 긴 흐름. 한 단계 실패 시 이전 단계의 보상 (compensating) 트랜잭션 으로 논리적 롤백. [!IMP…
[Redis] Pub/Sub vs Streams: 휘발 신호 vs 영속 로그database-internals
정의 - Pub/Sub ( / ): 지금 듣고 있는 구독자 에게만 메시지가 전달되는 휘발성 신호. 영속 없음, ACK 없음. fan-out. - Streams ( / / / ):…

💬 댓글

사이트 검색 / 명령어

검색

스크롤 = 확대/축소 · 드래그 = 이동 · 0 = 원래 크기 · ESC = 닫기