Files
DongHyeonka d646c2f12f 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개가 남음.
2026-08-14 14:55:38 +09:00

3.7 KiB

마이그레이션 가이드

기존 Spring Kafka / Spring AMQP 코드에서

1. topic 이름을 코드에서 제거한다

// 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 성공 판정을 없앤다

// 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를 끈다

# Kafka
enable.auto.commit: false
# RabbitMQ
auto-ack: false

둘 다 validator가 강제로 거부한다. 타이머 기반 commit은 handler가 실행되기도 전에 메시지를 처리 완료로 표시한다.

4. handler에서 ack 호출을 제거한다

// 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 마이그레이션

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 (바뀌지 않는다 — 이것이 의도다)