PostgreSQL과 Go를 사용한 트랜잭셔널 아웃박스 패턴
데이터와 함께 이벤트를 기록하세요. 절대 분리하지 마세요.
동시에 성공해야 하는 두 개의 쓰기가 결국 각각의 실패로 이어집니다.
주문 서비스는 먼저 데이터베이스에 주문을 저장한 후, 메시지 브로커로 order.created 이벤트를 발행합니다.
이 두 작업은 순차적으로 실행됩니다.
그 사이에서 문제가 발생합니다: 브로커가 다운되거나 네트워크 타임아웃이 발생하거나, 프로세스가 재시작되거나 컨테이너가 퇴출됩니다. 데이터베이스 쓰기는 성공했지만, 발행은 실패했습니다. 새로운 주문에 대한 정보를 받아야 하는 하위 서비스는 이를 알지 못합니다. 고객이 전화를 하기 전까지는 아무도 이를 눈치채지 못했습니다.
이것이 바로 듀얼-쓰기 문제이며, 분산 시스템에서 가장 흔한 침묵의 데이터 손실 원인 중 하나입니다. 트랜잭셔널 아웃박스 패턴(Transaction Outbox Pattern)이 표준 해결책입니다.

듀얼-쓰기 문제(Dual-Write Problem)
이 실패 모드는 일단 눈에 띄면 이해하기 쉽습니다.
BEGIN;
INSERT INTO orders ... -- 성공
COMMIT;
PUBLISH order.created ... -- 실패, 크래시, 또는 도달하지 못함
데이터베이스와 메시지 브로커는 트랜잭션 경계를 공유하지 않습니다. 두 작업을 모두 되돌릴 수 있는 롤백이 없습니다. 순차적으로 저장 -> 발행을 수행하는 모든 서비스에는 이러한 구멍이 존재합니다. 이 패턴은 여러 형태로 나타납니다:
db.Save(order)뒤에events.Publish(OrderCreated{...})호출- 트랜잭션을 커밋한 후 외부 웹훅을 호출하는 HTTP 핸들러
- 하나의 큐에서 레코드를 처리하여 결과를 다른 곳에 쓰는 워커
모든 경우의 결과는 동일합니다. 한쪽은 성공하고 다른 쪽은 실패하며, 시스템은 모니터링이 불가능한 상태에 빠집니다. 개별 작업들이 어느 시점에는 성공한 것처럼 보였기 때문입니다.
재시도 루프는 이 문제를 해결하지 못합니다. 데이터베이스 커밋 후 발행을 재시도하는 것은 재시도 자체의 신뢰성이 보장될 때만 작동하지만, 이는 당신이 가지지 못한 정확히 한 번(exactly-once) 내구성 보장과 같은 것입니다.
트랜잭셔널 아웃박스 패턴의 동작 방식
아웃박스 패턴은 직접적인 발행을 제거함으로써 이 구멍을 없앱니다. 비즈니스 로직 내에서 브로커를 호출하는 대신, 비즈니스 데이터와 동일한 데이터베이스 트랜잭션 내에 outbox 테이블에 이벤트 레코드를 작성합니다. 별도의 백그라운드 프로세스(릴레이, relay)가 outbox 테이블을 읽어서 브로커로 발행합니다.
BEGIN;
INSERT INTO orders ... -- 비즈니스 데이터
INSERT INTO outbox_events ... -- 이벤트 레코드
COMMIT;
-- 릴레이 프로세스 (별도):
SELECT ... FROM outbox_events FOR UPDATE SKIP LOCKED;
PUBLISH order.created ...
UPDATE outbox_events SET processed_at = NOW() WHERE id = $1;
두 쓰기가 모두 성공하거나 모두 실패합니다. 이미 PostgreSQL으로부터 얻고 있는 트랜잭션 보장이 이제 이벤트 레코드에도 적용됩니다. 릴레이는 필요한 만큼 발행을 재시도할 수 있는데, 그 이유는 이벤트가 내구성 있는 저장소(outbox)에 있기 때문입니다. 릴레이가 비행 중(처리 중)에 크래시되면 재시작하여 재시도합니다. 최악의 경우 이벤트가 한 번 이상 발행되는 것인데, 이는 컨슈머를 멱등성(Idempotent)으로 설계함으로써 처리됩니다 (see Idempotency in Distributed Systems).
아웃스 테이블을 위한 PostgreSQL 스키마
스키마는 의도적으로 단순합니다:
CREATE TABLE outbox_events (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
aggregate_type VARCHAR(100) NOT NULL,
aggregate_id VARCHAR(100) NOT NULL,
event_type VARCHAR(100) NOT NULL,
payload JSONB NOT NULL,
attempts INT NOT NULL DEFAULT 0,
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
processed_at TIMESTAMPTZ
);
-- 부분 인덱스: 처리되지 않은 행만 인덱싱하여, 행이 처리될 때마다 작게 유지됩니다
CREATE INDEX idx_outbox_unprocessed
ON outbox_events (created_at)
WHERE processed_at IS NULL;
processed_at IS NULL 조건을 가진 created_at의 부분 인덱스는 중요합니다. 이를 사용하지 않으면 인덱스는 작성된 모든 이벤트를 포함하여 커지고, 릴레이의 폴링 쿼리는 시간이 지날수록 느려집니다. 이를 사용하면 인덱스는 대기 중인 행만 커버하므로, 총 발행된 이벤트 수가 얼마나 되든 안정 상태에서 인덱스는 작고 경계가 정해진 집합을 유지합니다.
주요 필드 선택 사항:
aggregate_type과aggregate_id는 이벤트가 속한 엔티티를 설명합니다. 순서 보장 및 라우팅에 유용합니다.event_type은 컨슈머가 기대하는 이벤트 이름입니다.payload JSONB는 이벤트 본체를 저장합니다. 필요시 쿼리할 수 있도록TEXT가 아닌JSONB를 사용합니다.attempts는 릴레이가 이 행을 발행하려고 시도한 횟수를 추적합니다. 재시도 제한 및 데드 레터 처리에 사용됩니다.processed_at은 대기 중인 행에 대해NULL이며, 릴레이가 성공적으로 발행하면 설정됩니다.
비즈니스 데이터와 아웃박스 이벤트를 하나의 트랜잭션으로 작성하기
비즈니스 로직은 단일 BeginTx / Commit 호출 내에서 두 레코드를 모두 작성합니다. 여기에는 발행 호출이 없습니다. 오직 데이터베이스 쓰기가 있을 뿐입니다.
type OrderService struct {
db *sql.DB
}
func (s *OrderService) CreateOrder(ctx context.Context, order Order) error {
tx, err := s.db.BeginTx(ctx, nil)
if err != nil {
return fmt.Errorf("begin tx: %w", err)
}
defer tx.Rollback()
if _, err := tx.ExecContext(ctx, `
INSERT INTO orders (id, customer_id, total, created_at)
VALUES ($1, $2, $3, NOW())
`, order.ID, order.CustomerID, order.Total); err != nil {
return fmt.Errorf("insert order: %w", err)
}
payload, err := json.Marshal(map[string]any{
"order_id": order.ID,
"customer_id": order.CustomerID,
"total": order.Total,
})
if err != nil {
return fmt.Errorf("marshal payload: %w", err)
}
if _, err := tx.ExecContext(ctx, `
INSERT INTO outbox_events
(aggregate_type, aggregate_id, event_type, payload)
VALUES ($1, $2, $3, $4)
`, "order", order.ID, "order.created", payload); err != nil {
return fmt.Errorf("insert outbox event: %w", err)
}
return tx.Commit()
}
tx.Commit()이 실패하면 주문 행과 아웃스 행 모두 영구 저장되지 않습니다. 성공하면 둘 다 데이터베이스에 존재함이 보장됩니다. 릴레이는 그 이후의 어느 시점(즉시, 1초 후, 또는 크래시 후 재시작 후)에든 이벤트를 발행할 수 있습니다.
비즈니스 레이어에서 필요한 유일한 코드 변경입니다. 패턴의 나머지는 릴레이에 구현됩니다.
Go 릴레이 구현체
릴레이는 타이머에 따라 outbox 테이블을 폴링하는 백그라운드 워커입니다. 처리되지 않은 행들의 버치를 가져와 각각을 발행하고 완료 표시를 합니다. 애플리케이션과 동일한 바이너리에서 실행하거나 별도 프로세스로 실행할 수 있습니다. 둘 다 가능하지만 동일한 바이너리가 운영 측면에서 더 단순합니다.
type OutboxRelay struct {
db *sql.DB
publisher Publisher
logger *slog.Logger
batchSize int
pollInterval time.Duration
maxAttempts int
}
func (r *OutboxRelay) Run(ctx context.Context) error {
ticker := time.NewTicker(r.pollInterval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return ctx.Err()
case <-ticker.C:
if err := r.processBatch(ctx); err != nil {
r.logger.Error("outbox relay batch failed", "err", err)
}
}
}
}
릴레이는 컨텍스트 취소(context cancellation)를 존중하므로, 그raceful shutdown과 통합하기가 매우 간단합니다. 컨텍스트의 수명과 취소 패턴에 대한 자세한 내용은 [Go context.Context Done Right](https://www.glukhov.org/ko/app-architecture/code-architecture/go-context-cancellation-timeouts/ “취소, 타임아웃, 요청 범위 값을 위한 Go 컨텍스트 마스터하기. HTTP 핸들러, 데이터베이스 호출, 백그라운드 워커, 고루틴 누수, 그리고 그raceful shutdown을 다룹니다.“를 참조하세요.
FOR UPDATE SKIP LOCKED: 동시성 워커 패턴
processBatch 함수는 동시 릴레이 워커를 안전하게 처리하기 위해 FOR UPDATE SKIP LOCKED를 사용합니다.
func (r *OutboxRelay) processBatch(ctx context.Context) error {
tx, err := r.db.BeginTx(ctx, nil)
if err != nil {
return fmt.Errorf("begin tx: %w", err)
}
defer tx.Rollback()
rows, err := tx.QueryContext(ctx, `
SELECT id, aggregate_type, aggregate_id, event_type, payload
FROM outbox_events
WHERE processed_at IS NULL
AND attempts < $1
ORDER BY created_at
LIMIT $2
FOR UPDATE SKIP LOCKED
`, r.maxAttempts, r.batchSize)
if err != nil {
return fmt.Errorf("query outbox: %w", err)
}
defer rows.Close()
type row struct {
id string
aggregateType string
aggregateID string
eventType string
payload json.RawMessage
}
var batch []row
for rows.Next() {
var e row
if err := rows.Scan(
&e.id, &e.aggregateType, &e.aggregateID, &e.eventType, &e.payload,
); err != nil {
return fmt.Errorf("scan row: %w", err)
}
batch = append(batch, e)
}
if err := rows.Err(); err != nil {
return err
}
for _, e := range batch {
if err := r.publisher.Publish(ctx, e.eventType, e.aggregateID, e.payload); err != nil {
r.logger.Error("publish failed", "event_id", e.id, "err", err)
if _, err := tx.ExecContext(ctx,
`UPDATE outbox_events SET attempts = attempts + 1 WHERE id = $1`, e.id,
); err != nil {
r.logger.Error("increment attempts failed", "event_id", e.id, "err", err)
}
continue
}
if _, err := tx.ExecContext(ctx,
`UPDATE outbox_events SET processed_at = NOW() WHERE id = $1`, e.id,
); err != nil {
return fmt.Errorf("mark processed: %w", err)
}
}
return tx.Commit()
}
FOR UPDATE SKIP LOCKED은 두 가지 일을 수행합니다. 첫째, FOR UPDATE는 트랜잭션이 끝날 때까지 선택된 행을 잠가 다른 트랜잭션이 이를 선택하지 못하도록 합니다. 둘째, SKIP LOCKED는 행이 이미 다른 트랜잭션에 의해 잠겨 있을 경우, 쿼리가 대기하는 대신 해당 행을 건너뛴다는 의미입니다. 그 결과, 여러 릴레이 워커가 병렬로 실행할 수 있으며 각각은 중복되지 않는 행의 하위 집합을 가져옵니다.
SKIP LOCKED가 없다면, 두 번째 워커는 첫 번째 트랜잭션이 커밋할 때까지 차단되며, 그 시점에는 이미 행들이 완료 처리되었을 것입니다. SKIP LOCKED를 사용하면 두 번째 워커는 기다리는 대신 즉시 다른 행들을 가져오므로, 안전한 수평 확장(horizontal scaling)이 가능합니다.
위 코드에서 스캔(scan)과 발행(publish)의 분리를 주목하세요: 모든 행은 발행 루프가 시작되기 전에 슬라이스(slice)로 스캔됩니다. 이는 브로커에 대한 네트워크 호출 사이에 열린 *sql.Rows 커서를 유지하여 트랜잭션을 필요 이상으로 길게 유지하는 것을 방지합니다.
멱등성과 중복 제거(Deduplication)
릴레이는 최소 한 번(at-least-once) 발행합니다. 이벤트 발행 후 processed_at 업데이트를 커밋하기 전에 크래시하면, 재시작 시 동일한 이벤트를 다시 발행합니다. 이는 피할 수 없는 것입니다. 분산 트랜잭션 조정자(Distributed Transaction Coordinator) 없이 데이터베이스와 메시지 브로커 간에 정확히 한 번(exactly-once) 전달을 구현하려면 이 트레이드오프가 필요합니다.
컨슈머는 멱등성(Idempotent)이어야 합니다. 가장 간단한 접근법은 processed_events 테이블에서 처리된 이벤트 ID를 추적하는 것입니다.
CREATE TABLE processed_events (
event_id UUID PRIMARY KEY,
processed_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
);
func (h *OrderHandler) HandleOrderCreated(ctx context.Context, eventID string, payload []byte) error {
// 이벤트 ID를 자연 키로 사용하여 중복 제거
_, err := h.db.ExecContext(ctx, `
INSERT INTO processed_events (event_id) VALUES ($1)
ON CONFLICT (event_id) DO NOTHING
`, eventID)
if err != nil {
return fmt.Errorf("dedup check: %w", err)
}
// 삽입이 실제로 발생했는지(1행) 또는 노오퍼(0행)였는지 확인합니다.
// 더 간단한 접근법: RETURNING을 사용하거나 영향을 받은 행 수를 확인합니다.
// 영향 받은 행이 0이면 중복이므로 건너뜁니다.
...
}
실무에서는 많은 팀이 브로커 자체의 중복 제거 헤더(예: 로그 컴팩션(Log-compacted) 토픽을 위한 Kafka의 key 필드, 또는 RabbitMQ의 message-id 헤더)를 reliance하고, 데이터베이스 수준의 중복 제거를 대체재로 사용합니다. 둘 다 적용 가능한 유효한 레이어입니다.
중복 제거 키로서 발행된 메시지 안에 아웃스 이벤트 id(UUID)를 포함하세요. 그러면 컨슈머는 선호하는 중복 제거 메커니즘에 관계없이 이를 사용할 수 있습니다.
재시도 정책과 독성 메시지(Poison Messages)
attempts 컬럼은 재시도 정책을 주도합니다. 릴레이는 attempts >= maxAttempts인 행을 건너뛰며, 해당 행을 데드 레터(Dead Letter)로 처리합니다. 별도의 프로세스나 운영자 경보가 이를 처리합니다.
단순한 데드 레터 뷰:
CREATE VIEW outbox_dead_letters AS
SELECT *
FROM outbox_events
WHERE attempts >= 5
AND processed_at IS NULL
ORDER BY created_at;
좋은 프로덕션 재시도 정책:
maxAttempts를 재시도의 비용에 따라 5~10으로 설정합니다.- 지수 백오프(Exponential backoff)를 고려하세요:
retry_after컬럼을 추가하고retry_after > NOW()인 행은 건너뜁니다. outbox_dead_letters의COUNT(*)이 임계값을 초과하면 경보를 설정합니다.- 수동 재시도 경로를 제공하세요: 특정 행에 대해
attempts = 0및retry_after = NULL로 재설정하는 관리자 엔드포인트 또는 스크립트입니다.
컨슈머의 버그나 스키마 불일치로 인해 일관되게 실패하는 독성 메시지(Poison messages)는 정상적인 메시지를 차단해서는 안 됩니다. 릴레이는 틱마다 버치를 처리하고 실패를 큐에서 제거하는 것이 아니라 시도 횟수를 증가시켜 표시하므로, 정상적인 행은 정상적으로 진행되고 독성 메시지는 시도가 쌓여 데드 레터 임계값에 도달할 때까지 대기합니다. 이 outbox_dead_letters 뷰는 브로커 네이티브 [데드 레터 큐](https://www.glukhov.org/ko/app-architecture/integration-patterns/dead-letter-queues/ “독성 메시지 포착, 재시도 대 폐기 정책 선택, 그리고 안전한 재생 전략 설계를 위한 실용적인 데드 레터 큐 가이드.“가 구현하는 동일한 패턴의 데이터베이스 측 버전입니다. – 임계값 이후 격리, 볼륨에 대한 경보, 그리고 재생 전에 의도적인 결정 요구.
이벤트 순서와 파티셔닝
폴링 쿼리는 created_at을 기준으로 정렬하므로, 버치 내에서 선입선출(FIFO) 순서를 제공합니다. 대부분의 사용 사례에는 이것이 충분합니다. 동일한 주문에 대해 order.updated가 order.created보다 먼저 발행되지 않도록 해야 하는 등 엄격한 엔티티별 순서가 중요한 경우, 애그리게이트별 순서(Aggregate ordering)가 필요합니다.
ORDER BY 절에 aggregate_id를 추가하고 파티셔닝된 토픽([Apache Kafka](https://www.glukhov.org/ko/data-infrastructure/stream-processing/apache-kafka/ “타르볼 또는 Docker로 Apache Kafka 4.2 빠르게 학습하기, 로컬 KRaft 브로커 시작하기, 키 CLI 도구 마스터하기, 그리고 실용적인 프로듀서, 컨슈머, Connect 예제 실행하기.“에 발행할 때 메시지 키로 사용하세요. Kafka는 동일한 키를 가진 모든 메시지를 동일한 파티션으로 라우팅하며, 파티션은 순서대로 소비됩니다. 이는 전역 순서가 아닌 애그리게이트별 순서 보장을 제공합니다. 전역 순서는 단일 릴레이 인스턴스가 필요하기 때문입니다.
ORDER BY aggregate_id, created_at
파티셔닝된 순서를 지원하지 않는 브로커(예: 기본 AMQP 큐)의 경우, 단일 인스턴스 릴레이 또는 컨슈머 측 애플리케이션 레벨 순서 확인이 실용적인 대안이 됩니다.
LISTEN/NOTIFY를 사용하여 폴링 지연 시간 줄이기
1초의 폴링 간기는 평균 이벤트 지연 시간 500ms를 의미합니다. 대부분의 워크로드에는 괜찮습니다. 거의 제로(Zero) 지연 시간이 필요한 경우, PostgreSQL의 LISTEN/NOTIFY 메커니즘을 사용하면 새로운 아웃스 행이 삽입될 때 릴레이가 즉시 깨어날 수 있습니다.
아웃스 테이블에 트리거를 추가하세요:
CREATE OR REPLACE FUNCTION notify_outbox_insert() RETURNS trigger AS $$
BEGIN
PERFORM pg_notify('outbox_event', NEW.id::text);
RETURN NEW;
END;
$$ LANGUAGE plpgsql;
CREATE TRIGGER outbox_insert_notify
AFTER INSERT ON outbox_events
FOR EACH ROW EXECUTE FUNCTION notify_outbox_insert();
릴레이에서는 채널에서 리스닝하고 알림이 발생하면 깨어나되, 주기적인 폴링으로 대체(fallback)되도록 합니다.
func (r *OutboxRelay) Run(ctx context.Context) error {
listener := pq.NewListener(r.dsn, 10*time.Second, time.Minute, nil)
defer listener.Close()
if err := listener.Listen("outbox_event"); err != nil {
return fmt.Errorf("listen: %w", err)
}
ticker := time.NewTicker(5 * time.Second) // 폴링 대체재
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return ctx.Err()
case <-listener.Notify:
if err := r.processBatch(ctx); err != nil {
r.logger.Error("outbox batch failed (notify)", "err", err)
}
case <-ticker.C:
if err := r.processBatch(ctx); err != nil {
r.logger.Error("outbox batch failed (poll)", "err", err)
}
}
}
}
대체 타이커(Ticker)는 릴레이 재시작 중 또는 네트워크 문제 발생 시 놓친 알림을 처리합니다. 대체 간기는 밀리초가 아닌 몇 초로 유지하세요. 그 역할은 저지연이 아니라 복구입니다.
관측 가능성: 메트릭, 로그, 그리고 경보
아웃스는 인프라입니다. 인프라처럼 취급하고 그에 맞게instrumentation하세요.
주요 메트릭:
var (
outboxPublished = prometheus.NewCounter(prometheus.CounterOpts{
Name: "outbox_events_published_total",
Help: "성공적으로 발행된 아웃스 이벤트 총 수.",
})
outboxFailed = prometheus.NewCounterVec(prometheus.CounterOpts{
Name: "outbox_events_failed_total",
Help: "이벤트 유형별 아웃스 발행 실패 총 수.",
}, []string{"event_type"})
outboxPending = prometheus.NewGauge(prometheus.GaugeOpts{
Name: "outbox_events_pending",
Help: "현재 처리되지 않은 아웃스 이벤트 수.",
})
outboxBatchDuration = prometheus.NewHistogram(prometheus.HistogramOpts{
Name: "outbox_batch_duration_seconds",
Help: "각 아웃스 처리 버치의 지속 시간.",
Buckets: prometheus.DefBuckets,
})
)
가auge 갱신: outbox_events_pending의 정확성을 유지하기 위해 주기적인 쿼리를 실행하세요:
SELECT COUNT(*) FROM outbox_events WHERE processed_at IS NULL;
고려할 경보 임계값:
outbox_events_pending > 1000이 2분 이상 지속됨: 릴레이가 뒤처지거나 멈춤.outbox_events_pending이 단조 증가 중: 브로커 다운 또는 릴레이 크래시.- 데드 레터 카운터가 0이 아님: 스키마 또는 컨슈머 버그 조사 필요.
outbox_batch_duration_seconds p95 > 5s: 데이터베이스 지연 또는 버치 크기 과다.
구조화된 로그 필드: 릴레이의 모든 로그 라인에 event_id, event_type, aggregate_id, attempt를 포함하세요. 이 필드들은 실패한 발행을 특정 아웃스 행과 하위 컨슈머 추적과 상관관계 있게 연결할 수 있게 해줍니다.
아웃스 vs 직접 큐 vs 사가 패턴
아웃스 패턴은 모든 조정 문제에 적합한 도구는 아닙니다. 여기 비교가 있습니다:
| 접근 방식 | 원자성(Atomicity) | 복잡도 | 사용 시기 |
|---|---|---|---|
| 직접 발행 | 없음 | 낮음 | 이벤트 유실을 가끔 허용할 수 있는 경우 |
| 트랜잭셔널 아웃스 | 강력함 | 중간 | 단일 서비스로부터의 신뢰할 수 있는 이벤트 전달 |
| 사가 패턴 | 궁극적(Eventually) | 높음 | 여러 데이터베이스에 걸쳐 있는 다중 서비스 트랜잭션 |
| 2상 커밋(2PC) | 강력함 | 매우 높음 | 거의 실용적이지 않음; 대부분의 분산 시스템에서 피함 |
아웃스 패턴은 단일 서비스가 자신의 상태 변경을 반영하는 이벤트를 신뢰할 수 있게 방출함을 보장합니다. 여러 서비스에 걸쳐 상태 변경을 조정하는 것은 아닙니다. 그것이 [사가 패턴](https://www.glukhov.org/ko/app-architecture/integration-patterns/saga-transactions-in-microservices/ “마이크로서비스를 위한 Go 구현 사가 패턴 완전 가이드. 오케스트레이션 대 코오디네이션, 보상 전략, 멱등성, 이벤트 소싱, 그리고 실용적인 예제와 함께 모범 사례를 배웁니다.“의 역할입니다. 브로커 선택(예: RabbitMQ, SQS, 또는 Kafka)은 아웃스 패턴 자체와 독립적입니다. 릴레이는 시스템이 사용하는 브로커로 발행합니다.
사가를 구축 중이라면, 아웃스 패턴은 여전히 유용합니다. 사가의 각 참여자는 아웃스를 사용하여 로컬 상태 변경과 사가 이벤트를 하나의 트랜잭션으로 작성한 후, 사가 오케스트레이터나 코오디네이션이 이러한 이벤트를 신뢰할 수 있게 읽습니다.
WAL 기반 CDC를 대체 릴레이로 사용하기
폴링하는 대신 PostgreSQL의 Write-Ahead Log(WAL)를 테일링(tailing)하여 복제 스트림에서 아웃스 삽입을 직접 읽을 수 있습니다. Debezium과 같은 도구들이 이를 수행합니다. 이점으로는 더 낮은 지연 시간과 아웃스 테이블에 대한 잠금 압력(lock pressure) 감소가 있습니다. 단점으로는 운영 복잡도, 전용 PostgreSQL 복제 슬롯, 그리고 실행 및 모니터링을 위한 외부 서비스가 필요하다는 점이 있습니다.
대부분의 팀에게는 위에서 설명한 폴링 릴레이가 적절한 시작점입니다. 높은 아웃스 삽입률(초당 수만 건), 100ms 미만의 이벤트 지연 시간이 필요하거나, 다른 변경 캡처(BDC) 요구사항을 위해 이미 Debezium을 실행 중인 경우 WAL 테일링이 적합합니다.
sqlc 통합
sqlc를 사용하여 타입 안전한 Go 데이터베이스 코드를 작성한다면, 아웃스 쿼리는 자연스럽게 맞아떨어집니다:
-- name: InsertOutboxEvent :exec
INSERT INTO outbox_events (aggregate_type, aggregate_id, event_type, payload)
VALUES (@aggregate_type, @aggregate_id, @event_type, @payload);
-- name: FetchOutboxBatch :many
SELECT id, aggregate_type, aggregate_id, event_type, payload
FROM outbox_events
WHERE processed_at IS NULL
AND attempts < @max_attempts
ORDER BY created_at
LIMIT @batch_size
FOR UPDATE SKIP LOCKED;
-- name: MarkOutboxProcessed :exec
UPDATE outbox_events SET processed_at = NOW() WHERE id = @id;
-- name: IncrementOutboxAttempts :exec
UPDATE outbox_events SET attempts = attempts + 1 WHERE id = @id;
-- name: OutboxPendingCount :one
SELECT COUNT(*) FROM outbox_events WHERE processed_at IS NULL;
sqlc는 각 쿼리에 대해 타입 안전한 함수를 생성하므로, 문자열 보간 오류를 피하고 아웃스 쿼리 로직을 데이터베이스 액세스 레이어의 나머지 부분과 함께 유지할 수 있습니다.
프로덕션 체크리스트
아웃스 구현을 출시하기 전에 다음을 사용하세요:
데이터베이스
- 아웃스 테이블에
created_at WHERE processed_at IS NULL에 대한 부분 인덱스가 있음 - 기본값이 0인
attempts컬럼이 있음 - 데드 레터 뷰 또는 쿼리가 정의됨
- 오래된 처리된 행은 주기적으로 아카이빙되거나 삭제됨(야간 정리Job으로 충분함)
릴레이
- 폴링 쿼리에
FOR UPDATE SKIP LOCKED사용 - 릴레이가 트랜잭션 내에서 실행됨(쿼리 전에 시작, 모든 업데이트 후 커밋)
- 버치 크기가 제한됨(일반적으로 50~200행)
- 릴레이가 그raceful shutdown을 위해 컨텍스트 취소를 존중함
- 실패한 발행이 버치 중단이 아닌
attempts증가를 유발함
멱등성
- 발행된 메시지에 중복 제거 키로서 아웃스
id포함 - 컨슈머가 멱등성임 또는 브로커가 중복 제거를 제공함
- 중복 제거 패턴은 [Idempotency in Distributed Systems](https://www.glukhov.org/ko/app-architecture/integration-patterns/idempotency-in-distributed-systems/ “멱등성은 HTTP 트릭이 아닙니다. API, 큐, 웹훅, 워크플로우 전반에 걸친 중복 쓰기와 재생 메시지, 이중 과금을 막는 방법을 배웁니다.” 참조
관측 가능성
-
outbox_events_pending가auge가 모니터링 및 경보됨 - 데드 레터 카운터가 경보됨
- 릴레이 버치 지속 시간이 추적됨
- 구조화된 로그에
event_id,event_type,aggregate_id포함
운영
- 데드 레터 행을 위한 수동 재시도 경로가 있음
- 릴레이 재시작 동작이 테스트됨(올바르게 재발행하는가?)
- 브로커 장애 동작이 테스트됨(아웃스가 올바르게 증가하고 소멸하는가?)
마지막으로
듀얼-쓰기 문제는 사고가 발생하기 전까지 한계 사례로 치부하기 쉽습니다. 트랜잭셔널 아웃스 패턴은 이미 가지고 있는 도구들(PostgreSQL 트랜잭션, 백그라운드 고루틴, 그리고 하나의 추가 테이블)로 이를 해결합니다. 릴레이는 구축, 운영, 그리고 추론이 모두 단순합니다.
대가는 컨슈머가 최소 한 번(at-least-once) 전달을 위해 설계되어야 한다는 것입니다. 이는 합리적인 트레이드오프입니다. 분산 트랜잭션 없이 데이터베이스와 브로커 간에 정확히 한 번(exactly-once) 전달을 실현하는 것은 실제로 불가능합니다. 그렇지 않다고 속이는 것은 실패 조건 하에서 이벤트를 침묵적으로 유실하거나 이중 처리하는 시스템으로 이어집니다.
데이터와 함께 이벤트를 작성하세요. 신뢰할 수 있게 중계하세요. 컨슈머를 멱등성으로 만드세요. 그것이 패턴 전체입니다.
이 기사는 [App Architecture in Production](https://www.glukhov.org/ko/app-architecture/ “Slack 및 Discord 기반 통합 패턴, Python 클린 아키텍처 디자인 패턴, 그리고 GORM, Ent, Bun, sqlc 전반의 Go 데이터 액세스 트레이드오프를 포함하는 프로덕션 시스템을 위한 실용적인 앱 아키텍처 기둥.” 클러스터의 일부입니다.