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개가 남음.
114 lines
3.7 KiB
Markdown
114 lines
3.7 KiB
Markdown
# 마이그레이션 가이드
|
|
|
|
## 기존 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` (바뀌지 않는다 — 이것이 의도다)
|