feat(messaging): 브로커 중립 메시징 플랫폼 24개 leaf 추가

messaging-superpowers-package 설계서/계획서 기반 구현.
registry를 19 → 43 leaf로 확장하고 src/messaging 아래 24개 leaf를 등록.

- core-api: M1 publish/consume + M2 batch·delayed·pause-resume
- policy/transport-spi: 재시도 결정, DLQ orchestration, admission control, lifecycle
- kafka·rabbit(Stable): contiguous commit, confirm/return 상관, 배치, 보안 설정
- pulsar·nats(Experimental): 기본 비활성, live 인증 없음을 코드로 기록
- outbox/inbox/claim-check: 트랜잭션 결합, lease, 무결성 검증
- admin: plan → approve → execute를 타입으로 강제
- 문서 9종, infra compose 7종, JMH 벤치마크 3종

검증: 아키텍처 게이트 3종 통과, 24개 leaf 전부 check 통과,
messaging 테스트 604개 통과/0 실패.

미완: 계획서가 요구한 실 브로커 IT 40개 중 7개만 작성.
Rabbit 13 / Outbox 6 / Inbox 4 / NATS·Pulsar·Share 5 / testkit 2 /
starter·admin 3, 그리고 TLS·ACL 2개가 남음.
This commit is contained in:
DongHyeonka
2026-08-14 14:55:38 +09:00
parent 3b5aee50e3
commit d646c2f12f
486 changed files with 40271 additions and 2 deletions
+156
View File
@@ -0,0 +1,156 @@
# 설정 레퍼런스
## Destination profile
```yaml
messaging:
destinations:
order-events:
broker: kafka-primary
kind: EVENT_STREAM # ASYNC_COMMAND | DOMAIN_EVENT | INTEGRATION_EVENT
# | WORK_QUEUE | PUBLISH_SUBSCRIBE | EVENT_STREAM | REQUEST_REPLY
tier: M1 # M1 | M2 | M3
physical:
topic: order.events.v1
schema:
codec: application/json
compatibility: BACKWARD_TRANSITIVE
message-types: [order.created]
guarantees:
delivery: AT_LEAST_ONCE # AT_MOST_ONCE | AT_LEAST_ONCE
ordering: KEY # NONE | DESTINATION | PARTITION | KEY
external-side-effect: INBOX_TRANSACTIONAL
producer:
confirmation: REPLICATION_OR_PERSISTENCE_ACK
timeout: 5s
mandatory-routing: true
idempotent: true
consumer:
group: order-projection
concurrency: 6
max-in-flight-per-ordering-unit: 1
prefetch: 16
handler-timeout: 30s
manual-settlement: false
retry:
mode: PAUSE_PARTITION # NONE | INLINE | BLOCKING | PAUSE_PARTITION
# | RETRY_DESTINATION | BROKER_DELAYED
max-attempts: 3
initial-delay: 200ms
max-delay: 2s
multiplier: 2.0
jitter: true
ordering-impact: PRESERVE # PRESERVE | ALLOW_REORDER
dlq:
destination: order-events-dlq
max-redrive-count: 1
payload:
max-bytes: 1048576
claim-check-threshold-bytes: 1048576
key-resolver-configured: true
production: true
topology-auto-create: false
```
## 기본값
| 설정 | 기본값 | 근거 |
|---|---:|---|
| logical payload 최대 | 1,048,576 bytes | portability. 초과는 Claim Check |
| global hard 최대 | 8,388,608 bytes | 어떤 destination도 넘을 수 없는 상한 |
| header 총 크기 | 32,768 bytes | |
| header 개수 | 64 | |
| header key | 128 bytes | metric tag 안전 |
| header value | 4,096 bytes | |
| publish timeout | 5s | |
| handler timeout | 30s | |
| graceful shutdown drain | 30s | |
| 일반 destination retry | 0회 | 자동 retry는 opt-in |
| DLQ redrive batch | 100 | 한 번의 작업이 source를 덮치지 않게 |
| Outbox relay batch | 100 | |
| Outbox lease | 30s | |
| Outbox polling | 500ms | |
| metric dimension 상한 | 200 | cardinality 폭발 방지 |
## Broker profile
### Kafka
```yaml
messaging:
brokers:
kafka-primary:
type: kafka
stable: true
production: true
bootstrap-servers: [broker-1:9093, broker-2:9093]
enable-idempotence: true # stable에서 필수
acks: all # stable에서 필수
max-in-flight-requests-per-connection: 5 # 최대 5
delivery-timeout: 30s
enable-auto-commit: false # 항상 금지
tls-enabled: true # production 필수
authentication-enabled: true # production 필수
```
### RabbitMQ
```yaml
messaging:
brokers:
rabbit-primary:
type: rabbitmq
stable: true
production: true
addresses: [rabbit-1:5671]
publisher-confirms: true # stable에서 필수
publisher-returns: true # stable에서 필수
mandatory: true # stable에서 필수
confirm-timeout: 5s
auto-ack: false # 항상 금지
prefetch: 16
quorum-queues: true # durable work queue 필수
tls-enabled: true
authentication-enabled: true
```
## 보안
```yaml
messaging:
security:
kafka-primary:
producer: { type: SASL_SCRAM, credential-id: kafka-producer }
consumer: { type: SASL_SCRAM, credential-id: kafka-consumer }
# admin은 application runtime에 설정하지 않는다
hostname-verification: true
access:
publishable: [order-events]
consumable: []
administrable: []
```
## Experimental / Optional
기본값은 전부 `false`다.
```yaml
messaging:
experimental:
kafka-share: false
pulsar: false
nats: false
bridge:
spring-cloud-stream: false
```
## Backpressure
```yaml
messaging:
backpressure:
global-limit: 512
per-destination-limit: 64 # global-limit 이하여야 한다
```
`per-destination-limit > global-limit`이면 global limit이 limit이 아니게 되므로 부팅에 실패한다.
+75
View File
@@ -0,0 +1,75 @@
# 전달 보장
## 왜 `EXACTLY_ONCE`가 없는가
어떤 브로커도 **외부 side effect를 포함한** exactly-once를 제공하지 않는다.
실제로 존재하는 것은 at-least-once 전달 + 멱등하거나 transactional한 consumer의 조합이다.
플랫폼이 지킬 수 없는 이름을 enum에 두면 그 책임이 눈에 보이지 않는 곳으로 밀려난다.
그래서 `DeliveryGuarantee`는 증거가 끝나는 지점에서 멈춘다.
```java
public enum DeliveryGuarantee { AT_MOST_ONCE, AT_LEAST_ONCE }
```
## Publish 결과는 boolean이 아니다
```java
public enum PublishCompletion { CONFIRMED, REJECTED, AMBIGUOUS }
```
`REJECTED``AMBIGUOUS`를 하나의 "실패"로 합치면 중복 주문이 만들어진다.
전자는 broker가 저장하지 않았음이 **확정**되어 포기해도 안전하고, 후자는 그렇지 않다.
| 상황 | 결과 |
|---|---|
| 로컬 validation 실패 | `REJECTED`, `NOT_TRANSMITTED` |
| broker 명시적 reject / nack | `REJECTED` |
| confirm 수신 | `CONFIRMED` |
| Rabbit confirm + unroutable return | `REJECTED`, `UNROUTABLE` |
| bytes 전송 후 connection loss | `AMBIGUOUS` |
| confirm timeout | `AMBIGUOUS` |
| adapter가 판정 불가 | 보수적으로 `AMBIGUOUS` |
`PublishResult` 생성자가 이 규칙을 강제한다. `CONFIRMED`인데 broker acceptance가 없거나,
`AMBIGUOUS`인데 confirmation level을 주장하면 **객체 생성 자체가 실패**한다.
## Ordering
```java
public enum OrderingScope { NONE, DESTINATION, PARTITION, KEY }
```
순서는 partition·key·단일 consumer의 성질이지 destination 전체의 성질이 아니다.
`GLOBAL`이 없는 이유가 이것이다.
`DestinationProfileValidator`가 다음을 거부한다.
- `ordering=KEY`인데 key resolver 없음
- ordered destination인데 `ALLOW_REORDER` retry
- `orderingImpact=PRESERVE`인데 재발행형 retry(`RETRY_DESTINATION`, `BROKER_DELAYED`)
- `ordering=DESTINATION`인데 concurrency > 1
- ordered destination인데 ordering unit당 in-flight > 1
## External side effect
```java
public enum ExternalSideEffectGuarantee { NONE, IDEMPOTENCY_REQUIRED, INBOX_TRANSACTIONAL }
```
`INBOX_TRANSACTIONAL`만이 "DB side effect와 중복 차단이 같은 transaction에서 commit된다"를 의미한다.
Kafka transaction은 **Kafka 안에서만** 원자적이므로 이 값과 함께 설정하면
`KafkaTransactionProfileValidator`가 거부한다. 두 개의 독립적인 commit을 하나로 착각하게 두지 않기 위해서다.
## Consumer settlement 순서
```text
RECEIVED → DECODING → PROCESSING → HANDLER_SUCCEEDED → SETTLEMENT_SENDING
├→ SETTLED
└→ SETTLEMENT_UNKNOWN
```
- handler는 broker ACK API를 호출하지 않는다.
- `Success` 이후에만 source settlement한다.
- `SETTLEMENT_UNKNOWN`은 성공이 아니다. redelivery 가능성을 의미한다.
- `SettlementResult` 생성자가 `SETTLED`인데 `redeliveryPossible=true`인 조합을 거부한다.
+97
View File
@@ -0,0 +1,97 @@
# Experimental 정책
## Stable과 Experimental의 차이
**Stable**은 공통 Contract Suite(`MessagingAdapterContract`)를 변경 없이 통과한 어댑터다.
컴파일되는 어댑터가 아니라, 아래 7가지를 실제로 증명한 어댑터다.
```text
publishesAndConfirms
returnsAmbiguousWhenConfirmIsLost
redeliversWhenSettlementIsLost
preservesMessageIdAcrossRetryAndDlq
keepsSourceUnsettledWhenDlqPublishFails
rejectsOversizedPayloadBeforeTransport
stopsAcceptingNewWorkDuringShutdown
```
**Experimental**은 아직 그 증명이 끝나지 않은 어댑터다.
## 규칙
### 1. 기본 비활성
```yaml
messaging.experimental.kafka-share: false
messaging.experimental.pulsar: false
messaging.experimental.nats: false
```
활성화하지 않으면 validator가 `MessagingCapabilityUnavailableException`을 던진다.
Contract Suite가 아직 증명 중인 어댑터가 누군가의 기본 설정 때문에 load-bearing이 되어서는 안 된다.
### 2. Stable 모듈이 Experimental 모듈에 의존하지 않는다
Gradle 의존 그래프로 강제된다. `messaging-spring-boot-starter``allowed_dependencies`
`messaging-kafka-share-experimental`, `messaging-pulsar-experimental`,
`messaging-nats-experimental`, `messaging-spring-cloud-stream-bridge`**없다**.
`verifyCleanArchitectureDependencies`가 위반을 빌드 실패로 만든다.
### 3. Core 계약을 바꾸지 않는다
Experimental 어댑터는 브로커의 차이를 `MessagingCapabilities`로 표현할 뿐,
`messaging-core-api`의 타입을 바꾸지 않는다.
### 4. 없는 기능을 광고하지 않는다
| 어댑터 | 광고하지 않는 것 | 이유 |
|---|---|---|
| Kafka Share Group | orderedStream, keyedOrdering, replay, brokerTransaction | 경쟁 소비자 + 개별 ack는 partition 순서를 유지할 수 없다 |
| Pulsar | brokerTransaction | Pulsar에 있지만 플랫폼 Contract Suite로 증명되지 않았다 |
| Pulsar (Shared) | keyedOrdering | round-robin 분배 |
| NATS JetStream | nativeDeadLetter | delivery limit 초과 시 terminate할 뿐 라우팅하지 않는다 |
| NATS JetStream | keyedOrdering | subject 기반 모델에 per-key 순서가 없다 |
`false`인 capability를 요구하는 profile은 startup에서 실패한다.
조용히 약화되지 않는다.
### 5. 명시적 거부
| 조합 | 결과 |
|---|---|
| Kafka Share Group + ordering != NONE | 거부 |
| Kafka Share Group + pause/resume | `MessagingCapabilityUnavailableException` |
| Pulsar Shared + ordering=KEY | 거부 (Key_Shared 필요) |
| Pulsar + ordering=DESTINATION | 거부 |
| NATS Core + AT_LEAST_ONCE | 거부 (JetStream 필요) |
| NATS ordered consumer + 경쟁 워커 > 1 | 거부 |
| NATS + ordering=KEY | 거부 |
## Spring Cloud Stream bridge
Experimental이 아니라 **Optional**이다. 위험이 다르다.
Stream은 자체 binder 설정을 소유하므로, binding이 destination profile이 모르는
serializer·error handling·acknowledgement mode를 조용히 획득할 수 있다.
따라서 브리지는 **플랫폼 보장에 의존하지 않는 destination에만** 허용한다.
```text
ordering scope 선언 → 거부
retry policy 선언 → 거부
dead letter 선언 → 거부
```
이 셋 중 하나라도 필요하면 native adapter를 쓴다. 거기서만 실제로 강제되기 때문이다.
## 승격 조건
Experimental → Stable로 올리려면 전부 필요하다.
1. `MessagingAdapterContract` 7개 테스트를 변경 없이 통과
2. 장애 주입(연결 끊김, confirm 유실, settlement 유실) 하에서 통과
3. 지원 브로커 버전 범위 명시 및 CI 검증
4. `support-matrix.md`의 capability 표 갱신
5. ADR 작성
6. 기본 활성화 여부에 대한 별도 결정
+113
View File
@@ -0,0 +1,113 @@
# 마이그레이션 가이드
## 기존 Spring Kafka / Spring AMQP 코드에서
### 1. topic 이름을 코드에서 제거한다
```java
// before
kafkaTemplate.send("order.events.v1", key, payload);
// after
publisher.publish(orderEvents, envelope, PublishOptions.defaults());
```
`MessageDestination`은 logical name만 가진다. 물리 매핑은 destination profile이 소유한다.
`DestinationName`의 패턴이 `topic://orders` 같은 값을 거부하므로 우회할 수 없다.
### 2. boolean 성공 판정을 없앤다
```java
// before
try { template.send(...).get(); success(); }
catch (Exception e) { fail(); } // REJECTED와 AMBIGUOUS를 구분하지 못한다
// after
PublishResult result = ...;
switch (result.completion()) {
case CONFIRMED -> success();
case REJECTED -> abandon(); // broker가 저장하지 않음이 확정
case AMBIGUOUS -> retrySameMessageId(result); // broker가 가지고 있을 수 있음
}
```
이 구분이 없으면 confirm 유실 한 번이 중복 주문 하나가 된다.
### 3. auto-commit / auto-ack를 끈다
```yaml
# Kafka
enable.auto.commit: false
# RabbitMQ
auto-ack: false
```
둘 다 validator가 강제로 거부한다. 타이머 기반 commit은 handler가 실행되기도 전에
메시지를 처리 완료로 표시한다.
### 4. handler에서 ack 호출을 제거한다
```java
// before
@KafkaListener(...)
void handle(ConsumerRecord<?,?> record, Acknowledgment ack) {
process(record);
ack.acknowledge(); // 실패 시 순서가 애매해진다
}
// after
CompletionStage<HandleResult> handle(MessageDelivery<OrderCreated> delivery) {
process(delivery.message().payload());
return completedFuture(HandleResult.success());
}
```
settlement는 플랫폼이 수행한다. "성공한 뒤에만 ack"가 각 handler의 기억이 아니라
플랫폼 불변식이 된다.
### 5. 중복을 정상 상황으로 다룬다
at-least-once는 중복을 전제한다. 세 가지 중 하나를 고른다.
| 방식 | 언제 |
|---|---|
| handler 자체 멱등 | 자연 멱등 연산 (upsert 등) |
| Inbox | DB side effect가 있는 경우 |
| Kafka transaction | Kafka → Kafka 파이프라인만 |
`ExternalSideEffectGuarantee`에 선언한다. `INBOX_TRANSACTIONAL`과 Kafka transaction을
동시에 설정하면 거부된다. Kafka transaction은 DB를 포함하지 않는다.
### 6. 큰 payload는 Claim Check로
broker frame 크기를 키우지 않는다. broker 메모리, replication latency,
consumer recovery가 동시에 나빠지고, 유계·검증 가능한 실패가 무계 실패로 바뀐다.
1 MiB 초과는 외부 저장소로 offload하고 digest를 포함한 참조만 발행한다.
## DB 마이그레이션
```text
V1__messaging_outbox.sql
V2__messaging_inbox.sql
```
Outbox row는 business transaction과 같은 transaction에서 쓴다.
Inbox reservation은 handler side effect와 같은 transaction에서 쓴다.
별도 transaction이면 각 패턴이 닫으려던 창이 그대로 열려 있다.
## 단계적 전환
1. **publish만 전환** — 기존 consumer는 그대로. wire format은 reserved header가 추가될 뿐이다.
2. **Outbox 도입** — publish 유실 창을 닫는다.
3. **consume 전환** — handler를 `MessageHandler`로 옮기고 ack 호출을 제거한다.
4. **Inbox 도입** — 중복 side effect를 닫는다.
5. **retry·DLQ 정책 선언** — 이 시점까지 자동 retry는 0회다.
각 단계는 독립적으로 배포 가능하고, 되돌릴 수 있다.
## 되돌릴 수 없는 것
- 한 번 발행된 message type의 wire contract
- 이미 retention 안에 있는 메시지의 schema
- redrive된 메시지의 `messageId` (바뀌지 않는다 — 이것이 의도다)
+113
View File
@@ -0,0 +1,113 @@
# 운영 Runbook
## 배포 전 체크
```bash
./gradlew verifyCleanArchitectureDependencies --console=plain
./gradlew verifyRuntimeModuleMembership --console=plain
./gradlew verifyOneTypePerFile --console=plain
```
destination profile은 startup에서 검증된다. 아래는 **부팅 실패**다.
- ordered destination + reorder 가능 retry
- `ordering=KEY` + key resolver 없음
- payload 상한 > 8,388,608 bytes
- DLQ 자기 참조 / retry 자기 참조
- retry·DLQ 그래프 cycle
- 미등록 retry·DLQ destination
- M1 destination + manual settlement
- `AT_LEAST_ONCE` + confirmation `NONE`
- production profile + topology auto-create
- broker topology가 manifest와 불일치
## 증상별 대응
### publish가 AMBIGUOUS로 쏟아진다
broker confirm 경로 문제다. 실패가 아니다.
1. `PublishEvidence.transmission``MAY_HAVE_BEEN_TRANSMITTED`인지 확인
2. Kafka: `delivery.timeout.ms`, ISR 상태, leader election 확인
3. Rabbit: confirm timeout, channel 상태 확인
4. Outbox를 쓰고 있다면 `status='AMBIGUOUS'` row가 같은 messageId로 재시도 중이다. **정상이다.**
5. consumer 쪽 Inbox가 중복을 흡수하는지 확인
`AMBIGUOUS`를 실패로 취급해 새 messageId로 재발행하지 말 것. 중복이 복구 불가능해진다.
### DLQ가 비어 있는데 메시지가 사라졌다
DLQ publish 실패 시 source는 settlement되지 않는다. 메시지는 source에 남아 재전달된다.
1. `msg.failure-code``DEAD_LETTER_*`인 로그 확인
2. DLQ destination이 실제로 존재하는지 (topology validation)
3. DLQ credential에 publish 권한이 있는지
### consumer lag이 한 partition에서만 증가한다
`ContiguousPartitionOffsetTracker`가 gap에서 멈춘 것이다. 설계된 동작이다.
commit은 **연속** 완료 offset까지만 전진한다. offset 11이 아직 실행 중이면
10과 12가 끝나도 watermark는 10에 머문다. 12를 commit하면 consumer가 죽었을 때 11을 잃는다.
1. 해당 partition의 in-flight를 확인
2. 느린 handler를 찾는다 (`handlerTimeout` 초과 여부)
3. 필요하면 `PAUSE_PARTITION` retry가 걸려 있는지 확인
### 재시도 폭풍
`RetryPolicy.jitter=false`인지 확인한다. jitter 없이는 같은 초에 실패한 모든 consumer가
같은 초에 재시도한다.
### shutdown이 오래 걸린다
`GracefulShutdownCoordinator`가 in-flight를 기다리는 중이다.
- `inFlight()`가 0이 되면 즉시 종료
- drain deadline(기본 30초) 초과 시 남은 작업을 **unsettled로 포기**한다 → broker가 재전달
- draining 시작 후 새 retry attempt는 만들지 않는다
## Destructive 작업
전부 `DestructiveOperationGuard`를 통과해야 한다.
| 조건 | 요구 |
|---|---|
| admin credential | application runtime은 보유하지 않음 |
| `AdminApproval` | 유효기간 내 |
| dry-run | 항상 허용 |
### Replay
```text
기본: 격리된 consumer group (replay-<requestId>)
기존 group 대상: 승인 티켓 필수
```
기존 production group으로 replay하는 것은 "다시 읽기"가 아니라 **live consumer를 되감는 것**이다.
그 사이의 모든 것이 재처리된다.
### Redrive
```text
dry-run으로 후보 수 확인
→ 승인 획득
→ batch 100건 이하로 실행
→ republish CONFIRMED 인 것만 DLQ에서 settlement
```
`redriveId`로 재구동 루프를 추적한다. 같은 메시지가 반복해서 redrive되면
근본 원인이 해결되지 않은 것이다.
### Offset reset
`KafkaOffsetResetExecutor`는 승인 predicate를 **생성자 인자**로 받는다.
승인 소스 없이 조립된 runtime은 물리적으로 reset을 수행할 수 없다.
## Topology
production topology는 IaC가 만들고 애플리케이션은 **검증만** 한다.
`TopologyValidationRuntime`은 모든 불일치를 한 번에 보고하고 startup을 실패시킨다.
partition 수가 다르면 destination이 광고하는 ordering 보장이 달라지고,
`min.insync.replicas`가 없으면 `acks=all`의 의미가 달라진다.
+95
View File
@@ -0,0 +1,95 @@
# Outbox · Inbox
## 두 패턴이 각각 무엇을 해결하는가
| 패턴 | 해결하는 문제 | 해결하지 않는 문제 |
|---|---|---|
| Transactional Outbox | DB commit과 publish 사이의 창(窓) | 중복 |
| Inbox | 중복 delivery의 side effect | 유실 |
**둘 다 필요하다.** Outbox만으로는 exactly-once가 되지 않는다.
## Outbox
business transaction과 **같은 transaction**에서 row를 쓴다. 둘 다 commit되거나 둘 다 안 된다.
```sql
BEGIN;
UPDATE orders SET status = 'PLACED' WHERE id = ?;
INSERT INTO messaging_outbox (message_id, destination, ...) VALUES (?, ?, ...);
COMMIT;
```
### relay
```text
leaseBatch(100, 30s) -- lease로 다중 relay 인스턴스 안전
→ publish (messageId 그대로)
→ CONFIRMED → markPublished
→ AMBIGUOUS → markAmbiguous (같은 messageId로 재시도 가능)
→ REJECTED → markFailed
```
### 핵심 규칙: ambiguous는 같은 messageId로 재시도
새 id를 발급하면 "전달됐을 수도 있는 메시지"가 "확실히 두 번째인 메시지"가 되어
downstream의 어떤 중복 제거도 복구할 수 없다.
failed로 표시하면 broker가 이미 가지고 있을 수 있는 메시지를 잃는다.
`message_id`를 primary key로 둔 것도 같은 이유다. 어떤 코드 경로도 실수로 새 id를 붙일 수 없다.
### lease
```text
status IN ('PENDING','AMBIGUOUS') AND (lease_expires_at IS NULL OR lease_expires_at <= now)
```
partial index `ix_messaging_outbox_claimable`이 이 쿼리를 backlog 크기에 비례하게 유지한다.
PUBLISHED row는 retention job이 지울 때까지 쌓이기 때문이다.
## Inbox
reservation과 side effect가 **같은 transaction**이어야 한다.
```java
transactions.inTransaction(() -> {
if (!inbox.reserve(messageId, consumerId, now)) {
return InboxOutcome.duplicate(); // 이미 처리됨
}
return InboxOutcome.processed(sideEffect.get());
});
```
별도 transaction으로 예약하면 Inbox가 닫으려던 바로 그 창이 다시 열린다.
### 복합 키
`PRIMARY KEY (message_id, consumer_id)`.
message_id만으로 중복 제거하면 같은 event를 소비하는 두 번째 consumer가
첫 번째에 의해 억제된다. 각 consumer가 한 번씩 처리해야 한다.
### retention
broker의 최대 redelivery window보다 **길어야** 한다.
row를 먼저 지우면 늦게 도착한 redelivery가 두 번 처리된다.
## Debezium CDC 대안
polling relay 대신 WAL을 읽는다. polling interval과 lease 경합이 사라지지만
인프라와 그 자체의 실패 모드가 추가된다.
wire contract는 동일하다. `DebeziumOutboxEventRouter`가 polling relay와 같은 reserved header를
방출하므로 consumer는 어느 쪽이 발행했는지 구분할 수 없고, 전환은 배포 결정일 뿐 계약 변경이 아니다.
## Claim Check
1 MiB 초과 payload는 broker 프레임을 키우지 않고 외부 저장소로 offload한다.
`ClaimCheckReference`는 digest를 **필수**로 가진다. claim check는 메시지를 서로 다른 retention과
replication을 가진 두 시스템으로 쪼개므로, consumer는 producer가 저장한 바로 그 bytes를 받았음을
증명할 수 있어야 한다. 그렇지 않으면 잘린 객체와 정상 객체를 구분할 수 없다.
`ClaimCheckIntegrityGuard`는 fetch 전에 만료를, fetch 후에 크기와 digest를 검사한다.
digest 불일치는 `DESERIALIZATION`이 아니라 **validation** 실패로 분류한다.
bytes가 깨진 JSON인 게 아니라, 틀린 bytes이기 때문이다.
+103
View File
@@ -0,0 +1,103 @@
# Retry · DLQ · Redrive
## 자동 retry는 opt-in이다
일반 destination의 기본값은 **retry 없음**이다. 순서를 깨거나, 멱등하지 않은 side effect를
증폭시키거나, 이미 throttle된 downstream을 더 때리는 retry는 보이는 실패보다 나쁘다.
## 결정 순서
`DefaultRetryDecisionEngine`은 아래 순서를 위에서 아래로 평가한다.
```text
1. non-retryable category → parking(DeadLetter) 또는 Reject
2. attempt >= maxAttempts → DeadLetter
3. PRESERVE + ordered + orderedStream capability → PauseAndRetry
4. mode=PAUSE_PARTITION → PauseAndRetry
5. mode=RETRY_DESTINATION + ALLOW_REORDER → PublishToRetryDestination
6. mode=INLINE|BLOCKING → RetryInline
7. mode=BROKER_DELAYED + delayedDelivery capability → PublishToRetryDestination
8. 그 외 → DeadLetter
```
**retryability를 attempt 예산보다 먼저** 검사한다. deserialization 실패는 payload가 바뀌지 않으므로
재시도가 3번 더 실패할 뿐이다. 첫 delivery에서 바로 park한다.
**순서 보존 전략을 재발행 전략보다 먼저** 검사한다. 둘 다 설정되어 있어도 ordered destination이
reorder 경로로 흘러내리지 않는다.
## 기본 non-retryable
`DESERIALIZATION`, `AUTHENTICATION`, `AUTHORIZATION`, `CONFIGURATION`은 자동 retry하지 않는다.
매 redelivery마다 동일하게 실패하므로 부하만 늘어난다.
destination profile의 `retryableCategories`로 명시적으로 뒤집을 수는 있다.
## Backoff
`min(maxDelay, initialDelay * multiplier^(attempt-1))`, 이후 full jitter.
full jitter는 `[0, delay]` 균등 분포다. jitter가 없으면 같은 초에 실패한 모든 consumer가
같은 초에 재시도하고, downstream의 회복이 재시도 폭풍으로 즉시 무효화된다.
## Kafka: pause-and-seek vs retry topic
| 전략 | 순서 | 언제 |
|---|---|---|
| `PAUSE_PARTITION` | 유지 | ordered destination |
| `RETRY_DESTINATION` | 깨짐 | work queue, `ALLOW_REORDER` 명시 |
pause-and-seek는 메시지가 로그의 자기 자리를 떠나지 않는다. partition을 멈추고, 기다리고,
같은 offset으로 seek해 재전달한다. 뒤의 메시지도 함께 기다리며 이것이 의도된 동작이다.
## RabbitMQ: delayed retry queue
core broker에 per-message delay가 없으므로 **TTL + DLX**로 구현한다.
retry queue의 `x-message-ttl`이 만료되면 `x-dead-letter-exchange`를 통해 work queue로 되돌아간다.
주의: TTL 만료는 큐 **head**에서 평가된다. 하나의 retry queue에 서로 다른 delay를 섞으면
독립적으로 만료되지 않는다.
`basic.nack(requeue=true)`는 사용하지 않는다. delay 없이 큐 head로 되돌리므로 hot loop가 된다.
## DLQ: publish 확인 후 settlement
이것이 dead lettering이 데이터 손실이 되지 않게 하는 **유일한** 불변식이다.
```text
DLQ envelope 생성 (원래 messageId 유지)
→ DLQ publish
→ CONFIRMED 이면 source settlement
→ REJECTED / AMBIGUOUS 이면 source를 settlement하지 않음
```
source를 먼저 ACK하면, DLQ publish가 실패했을 때 메시지의 사본이 **어디에도 남지 않는다**.
broker는 이미 해제했고 DLQ는 받지 못했다.
AMBIGUOUS DLQ publish는 중복을 만든다. 이것이 의도된 trade다. DLQ는 사람이 읽는 곳이고
중복은 알아볼 수 있지만, 손실은 복구할 수 없다.
## DLQ envelope 내용
reserved header에만 기록한다. payload에 넣지 않는다.
```text
msg.failure-category, msg.failure-code, msg.origin-destination,
msg.retry-attempt, msg.first-failure-at, msg.last-failure-at
```
stack trace, exception message, secret header, 실제 key는 **넣지 않는다**.
DLQ는 원본 topic보다 오래 보관되고 더 많은 사람이 읽는다.
## Redrive
M4 Admin 전용이다. `DestructiveOperationGuard`를 통과해야 한다.
- admin credential 필요 (application runtime은 보유하지 않는다)
- 유효기간 내 `AdminApproval` 필요
- dry-run은 항상 허용 (계획이 공짜여야 사람이 계획한다)
- batch 상한 100건
- source == target 금지
- `redriveId``messageId`와 별개다. 재구동 루프를 식별하기 위해서다.
redrive도 **publish → settlement** 순서다. republish가 confirm되지 않은 메시지는
DLQ에 남는다.
+92
View File
@@ -0,0 +1,92 @@
# Messaging 보안
## Credential 분리
producer / consumer / admin은 **서로 다른 credential**이다.
`MessageSecurityValidator`가 startup에서 강제한다.
```text
producer credential == consumer credential → 실패
admin credential == producer|consumer → 실패
production 프로필에 admin credential 존재 → 실패
```
마지막 규칙이 "애플리케이션은 topic을 purge할 수 없다"를 **구조적으로** 만든다.
runtime이 admin 자격 증명을 아예 보유하지 않으므로, 침해된 handler가 상승시킬 대상이 없다.
## Production 필수 조건
- TLS 활성
- TLS hostname verification 활성
- broker authentication 활성
- topology auto-create 비활성
Kafka는 추가로 `enable.idempotence=true`, `acks=all`,
`max.in.flight.requests.per.connection <= 5`, consumer auto-commit 금지.
RabbitMQ는 추가로 publisher confirm, publisher return, `mandatory=true`,
durable work queue의 quorum queue, consumer auto-ack 금지.
## Credential은 값이 아니라 참조다
`BrokerCredentialProfile`의 어떤 variant도 secret을 담지 않는다.
식별자만 보관하고 connect 시점에 `CredentialProvider`로 해석한다.
heap dump나 설정 출력에서 사용 가능한 credential이 나오지 않는다.
`CredentialIds``bearer `, `sk-`, `-----begin`, `eyJ` 같은 접두사를 거부한다.
참조가 들어갈 자리에 secret 자체를 붙여넣는 가장 흔한 사고를 막는다.
## Rotation
`CredentialRotationPlan.isDue()`는 만료 **전에** 참이 된다.
broker가 연결을 거부하기 시작한 시점에는 이미 publish가 실패하고 consumer가 멈춰 있다.
rotation은 세대 교체다. `DefaultMessagingRuntimeRegistry.install()`이 새 세대를 원자적으로
게시하고, 이전 세대는 마지막 lease가 닫힐 때까지 열려 있다가 닫힌다.
진행 중인 publish는 시작한 연결에서 confirm을 받는다.
drain deadline이 이 대기를 제한한다. 없으면 lease 하나가 새면 폐기된 credential이
무기한 열려 있고, rotation이 보안상 무의미해진다.
## Header
금지 header는 application·platform 양쪽에서 거부한다.
```text
Authorization, Proxy-Authorization, Cookie, Set-Cookie,
access_token, refresh_token, api_key, password, client_secret
```
credential이 header에 들어가면 broker storage, DLQ dump, 운영 도구에 남는다.
downstream redaction으로는 되돌릴 수 없다.
예약 header(`msg.*`, `traceparent`, `tracestate`, `baggage`)는 platform만 쓴다.
application이 `msg.id`를 설정할 수 있으면 Inbox 중복 제거와 DLQ 상관관계가 의존하는
logical identity가 호출자 제어가 된다.
## ACL
`DestinationAccessValidator`가 broker ACL **이전에** 검사한다.
broker ACL 거부는 애플리케이션 컨텍스트가 없는 연결 수준 오류로 도착하므로
"어느 모듈이 어디에 publish하려 했는가"가 조사 대상이 된다.
## 관측성 누출
`MessagingRedactor`는 denylist다.
- secret: authorization, cookie, token, password, secret, credential
- per-message identity: messageId, correlationId, causationId, partitionKey, key, offset, deliveryTag, sequence
- payload: payload, body, data
- 예외 상세: exceptionMessage, stackTrace
identity를 지우는 이유는 두 가지다. bounded metric을 message당 하나의 series로 만들고,
support log를 재식별 표면으로 만들기 때문이다.
`CardinalityGuard`는 dimension당 값 개수를 상한한다.
cardinality 사고는 점진적이지 않다. 테스트 10건에서는 멀쩡하고 운영에서 백엔드를 죽인다.
## 감사
`MessagingAuditEvent`는 replay, redrive, offset reset, purge, delete를 기록한다.
subject(운영자 identity), approval ticket, 그리고 redactor를 통과한 details만 담는다.
누가 무엇을 했는지 증명하되 payload의 두 번째 사본이 되지 않는다.
+116
View File
@@ -0,0 +1,116 @@
# Messaging 지원 매트릭스
플랫폼이 **무엇을 보장하는지**와 **무엇을 보장하지 않는지**를 브로커별로 고정한다.
여기 없는 조합은 지원되지 않는다.
## 브로커 등급
| 브로커 | 등급 | 인증 기준 | Stable 기능 | 제한 |
|---|---|---|---|---|
| Kafka | Stable | 4.2+ / 4.3.x | producer idempotence, consumer group, batch, pause/resume, replay, transaction capability | Share Group은 Experimental |
| RabbitMQ | Stable | 4.3.x | exchange/routing, publisher confirm, mandatory return, manual ACK, quorum queue, retry queue, DLQ | stream 및 특수 plugin 미지원 |
| Pulsar | Experimental | 4.0 LTS + 4.2 | typed publish/consume, Shared, Key_Shared, schema | transaction 미승격, 기본 비활성 |
| NATS JetStream | Experimental | 2.14.x | stream, durable consumer, explicit ACK, dedupe, replay | native DLQ 없음(플랫폼이 대행), 기본 비활성 |
| Artemis/JMS | Extension | 범위 밖 | adapter SPI만 | 별도 ADR + Contract Suite 통과 필요 |
## Capability 매트릭스
`MessagingCapabilities`가 런타임에 선언하는 값이다. `false`인 기능을 요구하는 destination profile은
**startup에서 실패**하며, 조용히 약화되지 않는다.
| Capability | Kafka | Kafka Share | RabbitMQ | Pulsar | NATS JS |
|---|---|---|---|---|---|
| brokerAcknowledgement | O | O | O | O | O |
| replicationOrPersistenceEvidence | O | O | O | O | O |
| perMessageSettlement | O | O | O | O | O |
| batchSettlement | O | X | X | O | O |
| orderedStream | O | **X** | X | X | O |
| keyedOrdering | O | **X** | X | Key_Shared만 | X |
| replay | O | X | X | O | O |
| delayedDelivery | X | X | retry queue로 대행 | O | X |
| brokerTransaction | O | X | X | 미승격 | X |
| deduplicatedPublish | O | X | X | X | O |
| nativeDeadLetter | X | X | O | O | **X** |
| topologyManagement | O | X | O | O | O |
Kafka Share Group이 ordering 전부 `X`인 것은 설계 결정이다. share group은 개별 record를
경쟁 소비자에게 나눠주고 개별 ack하므로 partition 순서를 유지할 수 없다. ordered destination을
share group에 설정하면 `KafkaShareProfileValidator`가 거부한다.
NATS JetStream의 `nativeDeadLetter=X`도 마찬가지다. JetStream은 delivery limit 초과 시 메시지를
**terminate**할 뿐 어디로도 라우팅하지 않으므로, 플랫폼이 DLQ publish를 직접 수행한다.
## 기능 등급
| 기능 | 등급 |
|---|---|
| Typed Publish·Consume | Stable M1 |
| At-least-once contract | Stable |
| Ambiguous publish 결과 | Stable |
| handler 성공 후 자동 settlement | Stable M1 |
| Batch / Manual settlement / Pause·Resume / Delayed / Replay 요청 | M2 |
| Broker transaction / partition / routing / subscription | M3 |
| Replay 실행 / Redrive / offset reset / purge / delete | M4 Admin |
| Kafka Share Group, Pulsar, NATS | Experimental |
| Spring Cloud Stream bridge | Optional |
## 무엇이 "Stable"을 증명하는가
Stable 등급은 두 가지를 **모두** 통과해야 한다. `CompatibilityMatrixTest`가 이 규칙을 강제한다.
### 1. 공유 Contract Suite (`MessagingAdapterContract`, 7개)
Kafka와 RabbitMQ가 동일한 7개 테스트를 변경 없이 통과한다. 결정적 하네스를 쓰므로
확인 유실·settlement 유실 같은 장애를 요청 시점에 재현할 수 있다.
### 2. 실 브로커 인증 (Testcontainers)
| 스위트 | 무엇을 증명하는가 |
|---|---|
| `KafkaBrokerIT` | `acks=all`이 실제 replication 증거를 만든다 / 잘못된 토픽은 `REJECTED` / 발행-소비 왕복에서 identity 보존 및 contiguous commit |
| `KafkaAmbiguityChaosIT` | 브로커를 `docker pause`로 멈춘 상태의 publish가 **`AMBIGUOUS`** 로 보고된다 (broker acceptance 없음, confirmation level `NONE`, 비-retryable) |
| `RabbitBrokerIT` | exchange가 confirm했는데 어떤 큐에도 바인딩되지 않은 publish가 **`REJECTED` + `UNROUTABLE`** 로 보고된다 |
| `OutboxPostgresIT` | 롤백된 트랜잭션은 발행 가능한 행을 남기지 않는다 / `SKIP LOCKED` lease가 두 relay를 분리한다 / ambiguous 행이 같은 `messageId`로 재클레임된다 |
| `InboxPostgresIT` | 재전달이 side effect를 두 번 적용하지 않는다 / 롤백은 예약도 되돌린다 |
Docker가 없으면 `DockerAvailability` 가드로 skip되며, 이 표의 항목은 그때 **검증되지 않은 것**으로 취급한다.
### 3. 장애 시나리오 커버리지 (`BrokerFailureMatrix`)
`NetworkFaultScenario`가 5개 시나리오와 **각각의 기대 결과**를 코드로 고정한다. 기대 결과를 어댑터별로
두지 않는 것이 핵심이다 — 어댑터마다 다른 답을 허용하면 공유 계약이 존재할 이유가 없다.
| 시나리오 | 시점 | 기대 결과 | 이유 |
|---|---|---|---|
| `connection-refused` | 전송 전 | `REJECTED` | 바이트가 나가지 않았으므로 broker가 가질 수 없다 |
| `connection-cut-after-write` | 전송 후 | `AMBIGUOUS` | broker가 저장했고 confirm만 유실됐을 수 있다 |
| `confirm-timeout` | 전송 후 | `AMBIGUOUS` | timeout은 부재의 증거가 아니라 증거의 부재다 |
| `settlement-lost` | settlement 중 | `REDELIVERED` | 미settlement 메시지는 재전달이 설계다 |
| `high-latency` | 전송 후 | `AMBIGUOUS` | 판단 시점에는 confirm 유실과 구별할 수 없다 |
`CrossBrokerContractSuite`가 릴리스 게이트로 이를 강제한다. Stable 어댑터는 5개 전부를 **실 브로커에서**
커버해야 하고, Experimental 어댑터는 `LIVE_BROKER` 커버리지를 주장할 수 없다. 커버리지는 *능력*이 아니라
*무엇을 실제로 돌렸는지*의 기록이다.
### 실 브로커가 실제로 잡아낸 결함
이 스위트들은 장식이 아니다. 작성 과정에서 결정적 테스트가 통과하는데 실 인프라에서 실패한
결함을 두 건 잡았다.
1. **Outbox `IN_FLIGHT` 고아 행** — lease 쿼리가 `PENDING`/`AMBIGUOUS`만 클레임 대상으로 봐서,
publish 도중 죽은 relay가 남긴 행이 lease 만료 후에도 영영 회수되지 않았다.
2. **Rabbit confirm 경합** — transport가 publish *후에* confirm을 등록해서, 연결 스레드에서
confirm이 먼저 도착하면 유실되고 호출자가 무한 대기했다.
둘 다 인메모리 double이 실제보다 관대해서 통과하고 있었다.
## 명시적 비지원
- 공통 `EXACTLY_ONCE` 설정 — `DeliveryGuarantee`에 상수가 존재하지 않는다.
- 전역 순서 — `OrderingScope``GLOBAL`이 존재하지 않는다.
- DB와 broker의 자동 원자 transaction, 기본 XA
- Java native serialization
- 무제한 payload·header, 무한 retry
- 운영 application에서의 topology 파괴 작업
- 일반 애플리케이션에 raw broker client 반환
- DLQ publish 확인 전 source ACK