생산자 재시도와 relay 중복을 견디도록 event_id를 모든 경계에 보냅니다.
분산 메시지 큐
설계
메시지 큐는 비동기 연결 장치가 아니라, 어떤 사건을 얼마나 오래 보관하고 누구에게 어떤 순서와 중복 가능성으로 전달할지를 정하는 운영 계약입니다.
30초 핵심 요약
토픽은 여러 파티션 로그로 나누고, replica는 확정된 데이터를 견디며, Consumer Group은 파티션마다 진행 위치를 관리합니다. 이 구조는 파티션 범위의 순서와 at-least-once 전달을 다루기 좋지만, 결제·이메일처럼 외부 업무 효과를 자동으로 정확히 한 번 만들지는 않습니다.
생산자와 소비자의 속도·장애를 분리한다
분산 메시지 큐의 첫 질문은 “어떤 제품을 쓸까”보다 “어떤 업무 객체의 순서가 필요한가, 중복과 유실을 어디까지 허용하는가”입니다. 주문·결제·색인·분석은 모두 비동기일 수 있지만 같은 보장을 요구하지 않습니다.
order_id처럼 순서를 원하는 업무 객체를 key로 정하고 hot key를 측정합니다.
결제 후속 처리와 분석이 서로의 lag·장애에 끌려가지 않게 나눕니다.
offset reset, retry, DLQ, replay의 범위와 승인 경로를 남깁니다.
복제 로그와 업무 효과 저장소를 분리한다
outbox는 업무 변경과 publish 요청 사이의 누락을 줄이고, broker cluster는 파티션별 로그와 replica를 유지합니다. Consumer Group은 처리 병렬성을 얻지만 외부 API 호출·DB 갱신은 별도의 멱등 경계가 필요합니다.
event_id와 도메인 고유 제약으로 보호하며, broker의 acknowledgement를 사용자 업무 성공과 같은 뜻으로 쓰지 않습니다.메시지 envelope와 업무 효과의 식별자를 분리한다
publish API의 수락, broker의 복제 acknowledgement, consumer의 업무 완료는 서로 다른 상태입니다. timeout 뒤 중복 발행될 수 있으므로 생산자와 relay는 같은 event_id를 유지하고, 소비자는 그 ID를 업무 멱등 키로 전달합니다.
| 데이터 | 주요 필드 | 제약·역할 |
|---|---|---|
| outbox_event | event_id, aggregate_id, type, payload, created_at | UNIQUE(event_id); 업무 변경과 함께 기록하는 발행 원장 |
| message envelope | event_id, topic, key, headers, schema_version | partition·offset과 연결되는 전달 메타데이터 |
| consumer_effect | consumer_name, event_id, effect_state | UNIQUE(consumer_name,event_id); 외부 효과 중복 방지 경계 |
| dead_letter | original_ref, failure_class, first_seen | 격리·분석·명시적 재처리의 근거 |
발행, 처리, 진행 위치를 독립적으로 확인한다
타임아웃은 실패가 아니라 결과를 모르는 상태일 수 있습니다. 따라서 publish와 consume 모두 “다시 시도했을 때 같은 업무 결과가 안전한가”를 먼저 설계합니다.
상태 변경과 event_id를 같은 원장에 기록해 발행 누락을 줄입니다.
relay가 topic, key, schema version을 정하고 acknowledgement를 기다립니다.
broker는 key가 정한 파티션의 로그와 replica 상태를 관리합니다.
consumer가 effect ledger 또는 외부 멱등 키로 중복을 먼저 막습니다.
effect 뒤 offset을 진행시키고, 실패·재시작은 재처리 규칙으로 설명합니다.
순서는 전역이 아니라 key가 정한 파티션 범위다
하나의 Consumer Group에서 파티션 하나는 정상적인 할당 상태에서 한 consumer가 맡습니다. 이는 객체별 순서를 해석하기 좋은 기반이지만, 여러 파티션·여러 topic·외부 effect를 합친 전역 총순서는 아닙니다.
| 선택 | 얻는 것 | 위험·검증 | 운영 신호 |
|---|---|---|---|
order_id key | 같은 주문의 상태 전이를 한 partition 순서로 읽기 쉬움 | 초대형 주문·retry가 hot partition을 만들 수 있음 | partition별 bytes, lag, key skew |
customer_id key | 사용자별 제한·순서 적용이 단순 | 한 고객의 burst가 다른 고객을 지연 | tenant별 처리량, oldest age |
| 무작위 key | 쓰기 분산과 병렬성이 큼 | 업무 객체 순서 보장 불가 | 재정렬 필요 비율 |
| partition 증설 | 처리량·consumer 병렬성 증가 | key mapping 변경, rebalance, 파일·연결 비용 | assignment 안정 시간 |
at-least-once 전달과 업무 exactly-once를 분리한다
consumer가 DB에 결제 후속 상태를 기록한 직후 process가 죽고 offset commit 전에 멈추면, 재시작 때 같은 메시지를 다시 읽을 수 있습니다. 이 중복은 broker의 결함이 아니라 at-least-once 경계에서 예상해야 하는 상황입니다.
UNIQUE(consumer_name,event_id) 또는 수신 API idempotency key로 효과를 한 번만 적용합니다.
보존은 재처리 창을 주지만, 느린 group이 만료 전에 따라오는지 lag·oldest age를 함께 봅니다.
retention_headroom_secondsquota, bounded retry, priority 분리로 폭주를 막고 poison message는 원인과 함께 격리합니다.
dlq_messages_total{reason}DLQ는 종료 상태가 아닙니다. 원본 참조, 오류 분류, 첫 발생 시각, 보관 만료, 재발행 승인자를 남기고 원인 수정 후 제한 범위로 replay합니다. offset을 무조건 앞으로 보내 backlog를 없애면 업무 효과 누락을 숨길 수 있습니다.
Retention, backpressure, DLQ를 복구 수명으로 운영한다
Retention은 consumer 완료 여부와 별개로 로그를 얼마나 남길지 정합니다. 느린 group이 보존 만료 전에 따라오지 못하면 재처리 기회를 잃으므로 lag만이 아니라 남은 retention headroom을 관찰해야 합니다.
재처리·감사에 유리하지만 저장·복제·개인정보 노출 비용이 함께 증가합니다.
tenant quota, publish 제한, bounded retry, priority queue로 무한 유입을 제어합니다.
poison message, 비호환 schema, 권한 오류를 원인·만료와 함께 격리합니다.
원인을 고친 뒤 topic·partition·시간 범위·consumer를 제한해 재처리합니다.
backlog가 커졌다고 consumer 수만 늘리면 DB·외부 API를 더 세게 압박하거나 rebalance를 자주 일으킬 수 있습니다. DLQ 전체를 이유 없이 일괄 재발행하는 것도 장애를 증폭시킬 수 있으므로 retry budget과 재처리 승인을 운영 계약으로 둡니다.
장애마다 유실·중복·지연을 따로 검증한다
leader election과 replica lag로 produce 실패·지연 또는 낮은 RPO 여유가 발생합니다.
같은 event가 재전달되어 이메일·외부 요청·DB 갱신이 중복될 수 있습니다.
deserialize 오류가 하나의 파티션 처리와 consumer lag를 막습니다.
결제·메일 API 지연이 consumer concurrency를 잡아먹고 retention 창을 위협합니다.
member churn과 heartbeat timeout이 assignment를 반복해 중단·중복 가능성을 키웁니다.
publish/consume이 멈추거나 넓은 ACL이 민감 topic 노출을 만들 수 있습니다.
이벤트 로그를 데이터 복제면으로 취급한다
이벤트는 broker, replica, consumer, archive, DLQ에 복제되고 재처리될 수 있습니다. 따라서 “내부 메시지”라고 본문·권한·감사를 느슨하게 두면 안 됩니다.
producer, consumer, 운영자를 분리하고 일상 처리에 shared admin을 쓰지 않습니다.
publish·consume·offset reset·DLQ replay 권한을 별도로 제한합니다.
비밀번호·token·원본 카드 정보 대신 안전한 reference와 좁은 조회를 사용합니다.
ACL 변경, replay, offset reset, 민감 topic 읽기를 시간·주체와 남깁니다.
lag만 보지 말고 내구성과 업무 효과를 함께 본다
produce latency, retry, under-replicated partitions와 leader election을 한 화면에서 봅니다.
messages_produced_totalgroup·topic·partition별 records lag, age, retention headroom을 분리합니다.
consumer_lag_secondseffect applied, duplicate suppressed, DLQ reason, replay result를 event_id로 연결합니다.
effect_applied_totalevent_id, topic, partition, offset, group, schema version을 넣고 payload 전체·인증 토큰은 제외합니다. 평균 lag가 정상이어도 한 hot partition의 오래된 레코드가 업무 SLO를 깨뜨릴 수 있습니다.복제·보존·재처리까지 비용으로 계산한다
| 비용 축 | 왜 늘어나는가 | 설계 대응 |
|---|---|---|
| 저장·복제 | 원본, replica, index, compacted segment와 장애 여유 | topic별 retention·archive·RPO를 분리하고 compression을 실측 |
| 네트워크·CPU | cross-zone replication, encryption, compression, fan-out | 평균이 아닌 peak bytes/s와 burst를 기준으로 headroom 산정 |
| 운영·복구 | schema 관리, DLQ 분석, replay 승인, on-call | 재처리 runbook·대시보드·audit을 초기부터 제품화 |
| 개인정보 위험 | 긴 retention과 여러 consumer가 payload 노출면을 넓힘 | 참조 ID, encryption, topic ACL, 명시된 삭제·보존 정책 |
전달 모델은 요구하는 복구 방식에 맞춘다
| 선택 | 장점 | 제약 | 적합한 경우 |
|---|---|---|---|
| log 기반 topic·partition | replay와 독립 fan-out에 유리 | offset·partition·rebalance 운영 필요 | 이벤트 스트림·분석·여러 독립 소비자 |
| AMQP queue·routing | 명시적 queue와 유연한 route | retention·replay는 제품별 계약 확인 필요 | 작업 분배·복잡한 라우팅 |
| managed queue | broker 운영 부담 감소 | 제한·비용·이식성·관측 범위 검토 | 작은 운영팀, 관리형 SLA 수용 |
| 동기 RPC | 짧은 request-response가 단순 | 생산자와 소비자 지연·장애가 결합 | 즉시 결과가 꼭 필요한 좁은 경로 |
요구 보장과 복구 경계를 순서대로 설명한다
“결제 완료 이벤트를 여러 후속 시스템으로 안전하게 전달하는 메시지 큐를 설계해 보세요.”
- 먼저 허용 가능한 유실·중복, 대상 객체별 순서, replay 기간과 외부 side effect를 확인한다.
- outbox·keyed publish·partition replica·독립 Consumer Group을 기본 경로로 설명한다.
- 파티션 순서의 범위와 hot key·rebalance의 운영 영향을 명시한다.
- effect ledger 또는 idempotency key로 업무 중복을 막고, commit 순서와 복구를 설명한다.
- lag·retention headroom·DLQ·replica·권한 오류와 복구 검증 지표를 덧붙인다.
구현체의 계약은 공식 문서에서 다시 확인한다
- Apache Kafka — Design — partitioned log, replication, consumer position의 공식 설계 문서.
- Apache Kafka — Producer Configs — acknowledgement와 idempotence 같은 producer 설정의 공식 문서.
- AMQP — Specifications — AMQP 사양과 버전 자료의 공식 진입점.
- RabbitMQ — AMQP 0-9-1 Model Explained — producer, exchange, queue, consumer, acknowledgement 모델의 공식 설명.
- NATS — JetStream Concepts — persistence, consumer, retention 관련 공식 개념 문서.
작성·검토·참고 자료
참고 자료
- Apache Kafka — Design: partitioned log, replication, consumer position과 전달 의미를 설명하는 공식 설계 문서.
- Apache Kafka — Producer Configs: acknowledgement와 idempotence 같은 producer 설정을 확인하는 공식 문서.
- AMQP — Specifications: AMQP 프로토콜 사양과 버전 자료의 공식 진입점.
- RabbitMQ — AMQP 0-9-1 Model Explained: producer, exchange, queue, consumer, acknowledgement 모델을 설명하는 공식 문서.
- NATS — JetStream Concepts: persistence, consumers, retention과 재전송 개념을 설명하는 공식 문서. 출처는 구현체와 프로토콜의 용어·계약을 확인하기 위한 1차 자료다. 특정 지역의 보존 의무, 개인정보 삭제, 결제·주문 업무의 재처리 승인 절차는 조직의 법무·보안·제품 정책과 실제 선택한 서비스의 최신 운영 문서로 별도 확정해야 한다.
사실 오류·출처 정정은 문의·정정 페이지로 알려 주세요.