513 lines
112 KiB
Markdown
513 lines
112 KiB
Markdown
---
|
||
title: branch / feature-kafka-consumer-inbox-contract
|
||
source_type: branch-note
|
||
status: raw
|
||
branch: feature-kafka-consumer-inbox-contract
|
||
parent_branch:
|
||
related_projects: [ca-skeleton]
|
||
governing_docs: [raw/project-notes/ca-skeleton-operational-contract]
|
||
tags: [branch, ca-skeleton, kafka, consumer, inbox, backpressure]
|
||
created: 2026-07-28
|
||
target_merge:
|
||
status_label: in-progress
|
||
id: BR-CA-SKELETON-OPERATIONAL-CONTRACT-064
|
||
kind: project-work-item
|
||
project: ca-skeleton-operational-contract
|
||
work_item: WI-CA-SKELETON-OPERATIONAL-CONTRACT-064
|
||
inherits: [DEC-CA-SKELETON-OPERATIONAL-CONTRACT-DELIVERY-SEMANTICS-001@1, DEC-CA-SKELETON-OPERATIONAL-CONTRACT-IDEMPOTENCY-OWNERSHIP-001@1]
|
||
refines: []
|
||
overrides: []
|
||
depends_on: [WI-CA-SKELETON-OPERATIONAL-CONTRACT-063, WI-CA-SKELETON-OPERATIONAL-CONTRACT-070]
|
||
imports: []
|
||
delegates: []
|
||
accepts_delegations: []
|
||
contract_packet: 1
|
||
---
|
||
|
||
# branch: feature-kafka-consumer-inbox-contract
|
||
|
||
> Layer: `raw/branch-notes/` — 단일 브랜치의 **TODO·결정·진행 기록**. 머지/종료 후 verified 결과는 `/ingest`로 `wiki/projects/`에 추출. 원본은 raw에 영구 보관.
|
||
> `status_label`: `in-progress` | `review` | `merged` | `abandoned`
|
||
> **스캐폴딩 상태** — 결정(D-row)·구현 가이드는 비어 있다. `/branch-spec feature-kafka-consumer-inbox-contract` 로 채운다.
|
||
|
||
<!-- section-id: branch-parent -->
|
||
## 부모 (필수)
|
||
|
||
- **Parent project (canonical SSOT)**: [[raw/project-notes/ca-skeleton-operational-contract]]
|
||
|
||
> 분해 근거: `docs/superpowers/specs/2026-07-28-ca-skeleton-production-capability-feature-decomposition-design.md` §4.2 — 기술 런타임 (Tier T). 본 branch 는 project §8.0 `WI-CA-SKELETON-OPERATIONAL-CONTRACT-064` 의 실행 단위다.
|
||
|
||
형제 branch (같은 부모의 다른 자식 — 인접 영역):
|
||
|
||
- [[raw/branch-notes/feature-kafka-producer-runtime-contract]]
|
||
- [[raw/branch-notes/feature-idempotency-ownership-protocol-contract]]
|
||
- [[raw/branch-notes/feature-background-job-async-contract]]
|
||
|
||
<!-- section-id: branch-contract-packet -->
|
||
## 브랜치 계약 패킷
|
||
|
||
> project Work Item 에서 내려온 실행 계약의 snapshot. 여기에는 **pinned pointer + 1줄 요약 + branch 적용점**만 쓰고 상세를 복제하지 않는다.
|
||
|
||
- **생성 시 프로젝트 개정**: `1`
|
||
- **패킷 스키마**: `contract_packet: 1`
|
||
- **완료 조건**: inbound leaf 등록·manual ack·rebalance·DLT·inbox 멱등 계약 test 가 통과한다
|
||
|
||
<!-- section-id: inherited-project-decisions -->
|
||
### 상속한 프로젝트 결정
|
||
|
||
| Decision Ref | Project Summary | Branch Application | Source |
|
||
|---|---|---|---|
|
||
| `DEC-CA-SKELETON-OPERATIONAL-CONTRACT-DELIVERY-SEMANTICS-001@1` | end-to-end 메시징 보증은 at-least-once 전달과 멱등 consumer·inbox로 표현하고 DB와 broker를 걸친 exactly-once를 주장하지 않는다 | consume 측 경계: 오프셋 ack 는 (비즈니스 write + inbox insert) DB 커밋 **이후**에만 수행하고(D3), 중복 재전달은 inbox dedupe(D10·D11)가 흡수한다. Kafka 트랜잭션으로 DB 를 포함한 exactly-once 를 주장하지 않는다 | [[raw/project-notes/ca-skeleton-operational-contract]] |
|
||
| `DEC-CA-SKELETON-OPERATIONAL-CONTRACT-IDEMPOTENCY-OWNERSHIP-001@1` | idempotency는 owner token 기반 claim·renew·complete·release 프로토콜을 쓰고 실행 lease와 replay TTL을 분리하며 보증 등급을 명시한다 | consume 측 경계: 기본은 insert-once inbox(D10)이며 owner token claim/renew 재사용은 **조건부**(worker fan-out 으로 zombie consumer 동시 처리가 가능해질 때 — D12). 프로토콜 자체의 owner 는 [[raw/branch-notes/feature-idempotency-ownership-protocol-contract]] | [[raw/project-notes/ca-skeleton-operational-contract]] |
|
||
|
||
<!-- section-id: branch-local-decisions -->
|
||
### 브랜치 지역 결정
|
||
|
||
> 상세 근거·선택 조건·Open Risk 는 아래 결정-근거 매핑 §의 동일 D-row 가 소유한다. 여기에는 요약과 relation 만 둔다(복제 금지).
|
||
|
||
| Decision ID | Decision | Relation | Supporting Claims | Status |
|
||
|---|---|---|---|---|
|
||
| D1 | 신규 inbound leaf `adapter:inbound:messaging-kafka` 를 module registry 에 등록(19→20)하고 의존은 기존 inbound leaf 4종과 동일하게 제한 | `local` | ca-tmpl `.harness/project/modules.yaml` (실측) | `proposed` |
|
||
| D2 | seam 은 유지하되 스켈레톤이 `spring-kafka` 기반 기본 구현을 제공하고, broker 미선택 기동에서는 Kafka auto-config 가 켜지지 않아야 한다 (#063 D2 와 같은 축) | `local` | [[raw/branch-notes/feature-kafka-producer-runtime-contract]] D2 | `needs-approval` (#063 D2 와 동시 승인) |
|
||
| D3 | `enable.auto.commit=false` + use case 성공과 DB 커밋 **이후에만** 오프셋 ack | `refines DEC-CA-SKELETON-OPERATIONAL-CONTRACT-DELIVERY-SEMANTICS-001@1` | `raw/official-docs/kafka-consumer-offset-commit-semantics-apache-javadoc.md#KAFKA-OFFSET-C2` | `proposed` |
|
||
| D4 | 순서 단위는 파티션 — 동일 파티션 레코드는 항상 직렬 처리(공유 단일 큐 금지) | `local` | `raw/company-tech-blogs/kafka-multi-tier-retry-topic-dlq-uber.md#UBER-REPROC-C2` | `needs-confirmation` |
|
||
| D5 | backpressure 는 큐 포화 시 `pause()`/drain 후 `resume()` — 레코드 거부(drop) 금지 | `local` | `raw/official-docs/spring-kafka-listener-container-pause-resume-backpressure.md#SPRK-PAUSE-C2` | `proposed` |
|
||
| D6 | `CooperativeStickyAssignor` + `max.poll.*` 명시 pin, `onPartitionsRevoked` 를 유일 커밋 체크포인트로 신뢰 금지 | `local` | `raw/official-docs/kafka-incremental-cooperative-rebalance-kip429.md#KIP429-C5` | `proposed` |
|
||
| D7 | 역직렬화 실패(poison)는 리스너 호출 이전 단계에서 감지하고 **non-retryable** 로 분류 | `local` | `raw/official-docs/spring-kafka-error-handling-deserializer-poison-record.md#SPRK-EHD-C3` | `proposed` |
|
||
| D8 | 재시도 기본은 **blocking bounded retry**(파티션 순서 보존), non-blocking retry topic 은 기본 기각 | `local` | `raw/official-docs/spring-kafka-non-blocking-retry-topic-ordering-loss.md#SPRK-RETRYTOPIC-C3` | `proposed` |
|
||
| D9 | dead-letter 회수 + 감사. 발행은 **application-core port 경유** — inbound leaf 가 producer 를 직접 보유하지 않음 | `local` | `raw/official-docs/spring-kafka-default-error-handler-dlt-fatal-exceptions.md#SPRK-ERRH-C4` | `proposed` |
|
||
| D10 | `InboxStorePort` 기록과 비즈니스 write 를 **동일 `TransactionPort.inWrite` 경계**에서 커밋 | `refines DEC-CA-SKELETON-OPERATIONAL-CONTRACT-IDEMPOTENCY-OWNERSHIP-001@1` | `raw/official-docs/idempotent-consumer-microservices-io.md#MSIO-IDEMPC-C3` | `proposed` |
|
||
| D11 | dedupe key = envelope `idempotencyKey` 1차 + `(idempotencyKey, eventType)` 복합 유니크. `(topic, partition, offset)` 단독 채택 금지 | `local` | [[raw/branch-notes/feature-domain-event-outbox-contract]] D12·D14 | `proposed` |
|
||
| D12 | owner token 프로토콜(#070) 재사용은 조건부 — worker fan-out 도입 시에만 | `local` | [[raw/branch-notes/feature-idempotency-ownership-protocol-contract]] §범위 (D-row 미확정) | `needs-approval` |
|
||
| D13 | inbox row TTL 은 **수치 미확정** — 관계식(`retention.ms`/replay 창 중 긴 쪽 이상 + 무한 보관 금지)만 고정 | `local` | (UNSUPPORTED_DECISION — 외부 근거 부재) | `proposed` |
|
||
| D14 | consumer/inbox 용 error code·metric·env key 는 registry 에 **없음** — 전부 "신규 제안" 으로만 표기 | `local` | ca-tmpl `docs/registries/*.yaml` (실측: 해당 row 부재) | `proposed` |
|
||
| D15 | 등록된 **`(topic, eventType)`** 조합만 소비하고 미등록 조합은 D7 경로로 회수. allowlist 는 코드 handler 등록부로 둔다 (`schemaVersion` 축은 envelope 확장 후) | `local` | `raw/official-docs/spring-kafka-default-error-handler-dlt-fatal-exceptions.md#SPRK-ERRH-C1` (라우팅 불일치 계열은 fatal) + outbox D12 (envelope `eventType`) | `proposed` |
|
||
|
||
<!-- section-id: declared-overrides -->
|
||
### 선언한 예외
|
||
|
||
| Override ID | Overrides | Reason | Approval | Status |
|
||
|---|---|---|---|---|
|
||
|
||
<!-- GENERATED: project-contract-imports:start -->
|
||
## 가져온 프로젝트 계약
|
||
|
||
| Ref | Owner | 요약 | Branch 적용 |
|
||
|---|---|---|---|
|
||
<!-- GENERATED: project-contract-imports:end -->
|
||
|
||
<!-- section-id: branch-goal -->
|
||
## 목표
|
||
|
||
- `WI-CA-SKELETON-OPERATIONAL-CONTRACT-064` 의 완료 조건을 구현한다: inbound leaf 등록·manual ack·rebalance·DLT·inbox 멱등 계약 test 가 통과한다
|
||
|
||
- 이슈:
|
||
- PR:
|
||
|
||
<!-- section-id: branch-scope -->
|
||
## 범위
|
||
|
||
### 포함 범위
|
||
|
||
- `adapter:inbound:messaging-kafka` leaf 신설과 모듈 registry migration
|
||
- application 성공 이후 manual acknowledgement
|
||
- handler/schema/version allowlist
|
||
- bounded concurrency·queue 와 pause/resume backpressure
|
||
- rebalance·`max.poll` 처리, poison/역직렬화 실패 분류
|
||
- retry topic 또는 지연 재시도, DLT 와 감사된 replay
|
||
- `InboxStorePort` scope 와 같은 트랜잭션 커밋 규칙
|
||
|
||
### 제외 범위
|
||
|
||
> 의도적으로 제외한 것. 면접 등에서 "이건 범위에 없었습니다"라고 답할 근거.
|
||
|
||
- producer 설정 — #063 소유
|
||
- idempotency owner token 프로토콜 자체 — #070 소유
|
||
- 심층 명세 — 본 노트는 스캐폴딩. D-row 와 §구현 가이드는 `/branch-spec` 이 채운다.
|
||
- project decision registry 변경 — owner 는 project-note
|
||
|
||
## 근거 (필수, 최소 1개+)
|
||
|
||
> 외부 근거 추가 수집 진행 중. `/branch-spec feature-kafka-consumer-inbox-contract` 단계에서 verbatim 인용과 함께 `raw/official-docs/` 로 추가 수집한 뒤 여기서 링크한다.
|
||
|
||
| Source | 정당화하는 결정 |
|
||
|---|---|
|
||
| [[raw/company-tech-blogs/kafka-multi-tier-retry-topic-dlq-uber]] | 다단계 retry topic + DLQ 대안이 실제 운영에서 어떤 전제(단계별 backoff, 순서 비보장 수용, idempotent consumer) 위에 성립하는지의 사례 근거 — `company-case-study` 등급, 공식 best practice 아님. 순서 비보장 전제가 ca-skeleton per-aggregate FIFO 계약과 충돌 가능한 지점을 표시 |
|
||
| [[raw/official-docs/kafka-consumer-offset-commit-semantics-apache-javadoc]] | ca-skeleton 의 Kafka consumer 가 "application use case 성공 + inbox/비즈니스 트랜잭션 커밋 이후에만 offset 을 커밋(ack)" 하는 계약을 채택하고 `enable.auto.commit` 자동 커밋을 기각하는 근거 — 자동 커밋의 at-least-once 전제조건(poll 이후 전량 소비 필요)과 수동 커밋의 중복 창(commit 직전 crash → 재소비) 메커니즘. **한계**: rebalance·backpressure·DLT 상세는 이 자료 범위 밖 |
|
||
| [[raw/company-tech-blogs/kafka-poison-pill-consumer-stuck-offset-confluent]] | 역직렬화 실패(poison pill)를 재시도 무의미로 분류하고 즉시 DLT 로 보내야 하는 근거 — poison pill 이 consumer offset 을 전진시키지 못한 채 무한 재시도 루프에 빠뜨리는 실패 메커니즘, 그리고 역직렬화 실패가 `poll()` 반환 이전에 발생해 리스너 레벨 예외 처리로는 잡을 수 없다는 것 — `company-case-study` 등급(Confluent 벤더 블로그), 공식 best practice 로 격상 금지 |
|
||
| [[raw/official-docs/spring-kafka-listener-container-pause-resume-backpressure]] | bounded queue 포화 시 consumer 를 그룹에서 이탈시키지 않고 소비만 멈추는 backpressure 를 `pause()`/`resume()` 로 구현하는 근거 — pause 중에도 `poll()` 이 계속되어 rebalance 를 회피한다는 공식 동작과 반영 시점(poll 경계 vs `pauseImmediate`). **한계**: 파티션 단위 pause API(`pausePartition`/`resumePartition`)는 이 자료 범위 밖(별도 페이지, 추가 수집 필요) |
|
||
| [[raw/official-docs/spring-kafka-pause-resume-partitions-on-listener-containers]] | D4(파티션별 독립 bounded queue)와 D5(포화 시 pause)를 **그 파티션만** pause 하는 형태로 구현할 근거 — `pausePartition(TopicPartition)`/`resumePartition(TopicPartition)` API(2.7~), poll() 경계 반영 시점, `isPartitionPauseRequested()`/`isPartitionPaused()` 상태 조회 API. **한계(중요)**: rebalance·재배정 시 파티션 pause 상태의 운명(유지/초기화)은 이 문서도, 컨테이너 레벨 자매 문서도 **다루지 않는다**(`SPRK-PAUSEPART-C5` — 부재 확인) — §구현 가이드 3 의 `onPartitionsAssigned` 행이 "pause 상태 반드시 초기화"를 `KIP429-C3`+`SPRK-PAUSE-C1` 근거로 적어 두었으나, 두 claim 모두 이 구체 동작을 직접 말하지 않아 `UNSUPPORTED_IMPL_DECISION` 재라벨 후보로 남는다(branch 소유자 판단 필요, 본 자료는 근거 부재만 보고) |
|
||
| [[raw/official-docs/spring-kafka-non-blocking-retry-topic-ordering-loss]] | non-blocking retry topic 체인의 채택/배제 선택 조건 근거 — 공식 문서가 자인하는 순서 보장 손실(SPRK-RETRYTOPIC-C3)을 per-aggregate 순서 보장 요구사항과 대조해 언제 이 대안을 배제하는지 판단하는 근거. 위 Uber 사례(`kafka-multi-tier-retry-topic-dlq-uber`)의 순서 비보장 전제가 공식 문서로도 뒷받침됨을 확인 |
|
||
| [[raw/company-tech-blogs/kafka-consumer-rebalance-cooperative-sticky-verygoodsecurity]] | `cooperative-sticky` 파티션 할당 전략 + `max.poll.*` 튜닝 병행 채택의 운영 사례 근거 — 잦은 rebalance 로 consumer 가 그룹에서 이탈하고 커밋이 실패하던 환경에서 두 조치를 적용한 결과. `company-case-study` 등급(VGS 벤더 블로그), 공식 best practice 로 격상 금지 — 규모(100 consumers/partitions, aiokafka Python 클라이언트) 전제가 ca-skeleton 과 다를 수 있음 |
|
||
| [[raw/official-docs/kafka-consumer-configs-max-poll-and-commit-defaults]] | `max.poll.interval.ms`(초과 시 그룹 이탈·rebalance)/`max.poll.records`/`enable.auto.commit`·`auto.commit.interval.ms`/`session.timeout.ms`·`heartbeat.interval.ms`/`partition.assignment.strategy` 각각의 공식 정의·기본값 기준선 — 임의 수치 발명 방지. 기본 `partition.assignment.strategy`(`[RangeAssignor, CooperativeStickyAssignor]`)가 이미 CooperativeStickyAssignor 로의 단일 rolling-bounce 업그레이드 경로를 지원함을 확인 |
|
||
| [[raw/official-docs/spring-kafka-default-error-handler-dlt-fatal-exceptions]] | poison/역직렬화 예외 6종 기본 fatal 분류(`SPRK-ERRH-C1`) · blocking retry backoff 가 consumer 스레드를 정지시켜 `max.poll.interval.ms` 초과 시 rebalance 위험을 만들고 그래서 `ContainerPausingBackOffHandler` 가 제공된다는 메커니즘(`SPRK-ERRH-C2`) · DLT 기본 명명 `<originalTopic>-dlt` + partition 요건 + recoverer 의 producer(`KafkaTemplate`/`KafkaOperations`) 요구(`SPRK-ERRH-C3`/`C4`) · recoverer 미구성 시 기본 동작이 로그만이라는 사실(`SPRK-ERRH-C5`) |
|
||
| [[raw/official-docs/kafka-incremental-cooperative-rebalance-kip429]] | `partition.assignment.strategy=CooperativeStickyAssignor`(incremental cooperative rebalance) 채택의 Kafka 공식 사양(KIP-429, Accepted 2.4.0) 근거 — EAGER 는 rebalance 마다 소유한 모든 파티션을 revoke 하지만 COOPERATIVE 는 소유 파티션을 유지한다는 정의(`KIP429-C2`/`KIP429-C3`), 그리고 cooperative 프로토콜에서 `onPartitionsRevoked` 가 아예 호출되지 않을 수 있어 이를 rebalance 시작 신호나 유일한 커밋 체크포인트로 신뢰하면 안 된다는 공식 근거(`KIP429-C5`) — VGS 사례(`kafka-consumer-rebalance-cooperative-sticky-verygoodsecurity`)의 `cooperative-sticky` 채택을 공식 사양으로 보강 |
|
||
| [[raw/official-docs/spring-kafka-error-handling-deserializer-poison-record]] | poison message(역직렬화 실패)를 리스너 호출 이전 단계(deserializer 레벨)에서 감지해 error handler/DLT 경로로 회수하는 방식 채택 — `ErrorHandlingDeserializer` 가 위임 deserializer 실패 시 null 값 + `DeserializationException` 헤더(원인 + raw bytes)를 실어 보내고, 컨테이너가 리스너 대신 `ErrorHandler` 를 호출한다는 Spring 공식 메커니즘(`SPRK-EHD-C2`/`C3`). Confluent 사례(`kafka-poison-pill-consumer-stuck-offset-confluent`)의 "재시도 무의미" 판단을 Spring 프레임워크 레벨의 구체적 반환값·라우팅 계약으로 보강 |
|
||
| [[raw/official-docs/idempotent-consumer-microservices-io]] | inbox(PROCESSED_MESSAGE류) 테이블에 처리한 메시지 ID 를 기록해 at-least-once 재전달 중복을 탐지·폐기하는 방식 채택 근거 — ID INSERT 가 message handler 의 DB 트랜잭션 경계 안에서 이뤄지고 (subscriberId, messageID) 복합 유니크 제약으로 duplicate INSERT 가 실패·rollback 된다는 메커니즘(`MSIO-IDEMPC-C3`/`C4`), 그리고 별도 테이블 대신 비즈니스 엔티티 자체에 ID 를 저장하는 변형 옵션(`MSIO-IDEMPC-C5`). `engineering-blog` 등급(Chris Richardson 개인 패턴 카탈로그) — 공식 벤더 문서로 격상 금지, "동일 트랜잭션 요구"의 명시적 문장은 미발견(원본 raw 의 Usage Boundaries 참고) |
|
||
| [[raw/official-docs/spring-kafka-ack-mode-manual-commit-and-concurrency]] | D3 의 manual ack 계약을 `ContainerProperties.AckMode` 층에서 어떻게 표현하는지의 근거 — `AckMode.MANUAL`/`MANUAL_IMMEDIATE` 정의(`SPRK-ACKMODE-C1`/`C2`), 기본값이 `MANUAL` 이 아니라 `BATCH` 라 명시 설정이 필요하다는 것(`SPRK-ACKMODE-C4`), `nack()`/`acknowledge(index)` 의 리스너·consumer 스레드 제약(`SPRK-ACKMODE-C5`/`C6`), `concurrency` > 파티션 수일 때 하향 조정(`SPRK-ACKMODE-C7` — D4/D6 의 "파티션당 컨슈머 1개" 전제와 정합, 단 Kafka 프로토콜 레벨 보장 자체의 대체 근거는 아님). **한계**: 사용자가 요청한 "Acknowledgment 를 별도 워커 스레드에서 호출해도 되는가"(일반 `acknowledge()` 의 스레드 규칙)·ack 순서 제약·`asyncAcks` trade-off 3가지는 이 페이지에서 self-grep 0건으로 미발견 확인 — 별도 페이지("Manually Committing Offsets") 조사 필요 |
|
||
|
||
**추가 수집 필요** (`/branch-spec` 단계): rebalance 시 `ConsumerRebalanceListener` 공식 API 세부 — **파티션 단위 pause/resume API 자체는 2026-07-28 `[[raw/official-docs/spring-kafka-pause-resume-partitions-on-listener-containers]]` 수집으로 해소되었으나, 그 문서도 rebalance·재배정 시 pause 상태의 운명은 다루지 않아 `ConsumerRebalanceListener`/`onPartitionsAssigned` 상호작용 근거는 여전히 미수집**, ca-tmpl platform 설계 §11.4, DLT 실제 라우팅 구성(`DefaultErrorHandler`+`DeadLetterPublishingRecoverer`)의 "Handling Exceptions" 공식 페이지
|
||
|
||
**프로젝트 내부 설계 참조 (등급 `internal-design-doc` — 공식 문서 아님, best practice 로 격상 금지):**
|
||
|
||
- ca-tmpl `docs/superpowers/specs/2026-07-26-production-capability-platform-design.md`
|
||
- llm-wiki `docs/superpowers/specs/2026-07-28-ca-skeleton-production-capability-feature-decomposition-design.md`
|
||
|
||
## 외부 근거 / 대안 조사 (2026-07-28 — `/branch-spec` 자동조사 6건)
|
||
|
||
결정마다 대안을 실제로 비교한 기록. 6개 결정 topic 을 `wiki-decision-researcher` 로 조사했고(회당 bound 6, 초과분 없음), 채택안·기각안의 근거 자료를 위 §근거 표에 raw 로 보존했다.
|
||
|
||
| # | 조사한 결정 topic | 비교한 대안 | 채택 / 기각 | 대응 D-row |
|
||
|---|---|---|---|---|
|
||
| 1 | 오프셋 커밋(ack) 방식 | (a) `enable.auto.commit=true` / (b) raw `commitSync`·`commitAsync` 수동 커밋 / (c) 컨테이너 ack 모드 위임 | (a) **기각** — "성공 후 ack" 목표를 구조적으로 만족 못 함. (b)/(c) 는 D2(SPI vs concrete)에 종속되는 구현 형태 차이 | D3, D2 |
|
||
| 2 | 장시간 처리와 rebalance 안정성 | (a) `max.poll.*` 튜닝 / (b) `pause()`+별도 워커 / (c) cooperative-sticky + rebalance listener | (b) 채택(범위상 기정) + (c) 병행 채택 + (a) 는 defense-in-depth 안전판 | D5, D6 |
|
||
| 3 | 동시성·backpressure | (a) 파티션당 전담 스레드 / (b) bounded queue + pause/resume / (c) reactive(Reactor Kafka) backpressure | (a)를 기본, 넘칠 때 (b) 로 확장(파티션별 독립 큐 강제). (c) **기각** — sibling 이 스레드 기반 어휘를 SSOT 로 확립해 어휘가 분기됨 | D4, D5 |
|
||
| 4 | poison·역직렬화 실패 | (a) deserializer 경계 감지 + 회수 / (b) `byte[]`/`String` 소비 후 application 파싱 / (c) skip-and-log | (a) 를 파싱 실패 경계로, (b) 를 도메인 규칙 위반 경계로 **역할 분담**. (c) **기각** — 감사 흔적 없이 조용히 유실 | D7 |
|
||
| 5 | 재시도 전략·DLT | (a) blocking retry + DLT / (b) non-blocking retry topic 체인 / (c) 외부 지연 큐·DB 기반 지연 재시도 | (a) 채택 — (b) 는 공식 문서가 순서 손실을 자인해 D4 와 충돌. (c) 는 dead-letter 를 DB row 로 두는 대안으로 D9 에 보존 | D8, D9 |
|
||
| 6 | consumer 중복 차단(inbox) | (a) inbox 테이블 + 동일 트랜잭션 / (b) 비즈니스 자연 유니크 제약 / (c) 외부 캐시(Redis) dedupe | (a) 채택(감사·범용성). (b) 는 1이벤트=1row 조건부 대안. (c) **기각** — Redis 는 optional adapter 라 정합성 근거를 mandatory 의존 밖에 두게 됨 | D10, D11, D12, D13 |
|
||
|
||
**비교의 핵심 축**: ① 순서 보장(파티션 단위)을 지킬 것인가 처리량을 살 것인가 — 5번이 여기서 갈린다. ② 정합성 근거를 mandatory 의존(PostgreSQL) 안에 둘 것인가 — 6번이 여기서 갈린다. ③ 스켈레톤이 SPI 인가 concrete 구현인가 — 1번의 (b)/(c) 선택이 여기에 종속되며 D2 가 이를 닫는다.
|
||
|
||
**조사 후에도 근거가 없어 라벨링한 것**: inbox TTL 수치(D13 `UNSUPPORTED_DECISION`), 감사된 replay 기록 스키마·계약 값 명명(§구현 가이드 4·6 의 `UNSUPPORTED_IMPL_DECISION`).
|
||
|
||
## TODO
|
||
|
||
각 항목 옆에 증거 등급 표기: `actually-implemented` | `locally-verified` | `prod-verified` | `documented-only` | `planned` | `needs-confirmation`
|
||
|
||
- [x] `/branch-spec feature-kafka-consumer-inbox-contract` 로 D-row·§구현 가이드 작성 — 등급: `documented-only` (D1~D15 + §구현 가이드 7절 작성, 2026-07-28)
|
||
- [ ] `SOURCE_GAP-1` 해소 — 파티션-소비자 배타 배정의 Kafka 공식 verbatim 수집 후 D4 격상. **이미 수집한 `raw/official-docs/kafka-consumer-offset-commit-semantics-apache-javadoc.md` 와 동일한 KafkaConsumer Javadoc 페이지의 §Consumer Groups and Topic Subscriptions 절**에 해당 문장이 있으므로 새 URL fetch 없이 claim 추가 추출로 닫힌다 — 등급: `planned`
|
||
- [x] `SOURCE_GAP-2` 해소(2026-07-28) — 파티션 단위 pause/resume API 공식 페이지(`[[raw/official-docs/spring-kafka-pause-resume-partitions-on-listener-containers]]`) 수집 완료: API 존재(`pausePartition`/`resumePartition`, 2.7~)·타이밍·상태조회는 확인. "재배정 시 pause 상태 초기화" 자체는 이 문서도 다루지 않음을 확인해 §구현 가이드 3 의 `onPartitionsAssigned` 행을 `UNSUPPORTED_IMPL_DECISION`(trade-off: 보수적으로 명시 resume)으로 강등 완료 — 등급: `documented-only`
|
||
- [x] `SOURCE_GAP-5` 대부분 해소(2026-07-28) — spring-kafka `AckMode` 공식 페이지(`[[raw/official-docs/spring-kafka-ack-mode-manual-commit-and-concurrency]]`) 수집: ack 모드 정의·**기본값 `BATCH`**·리스너 타입 제약·concurrency 하향 조정 확보 — 등급: `documented-only`
|
||
- [ ] `SOURCE_GAP-5` 잔여 — 일반 `acknowledge()` 의 호출 스레드 규칙과 ack 순서 제약은 위 페이지에 **부재 확인**. Spring Kafka "Manually Committing Offsets" 페이지에서 수집해 §구현 가이드 2 의 잔여 `UNSUPPORTED_IMPL_DECISION` 해소 — 등급: `planned`
|
||
- [ ] **D2 는 #063 D2 와 동시 승인** — 승인 전까지 Kafka SDK 를 classpath 에 반입하지 않는다(capability-provider D13 충돌 위험) — 등급: `planned`
|
||
- [ ] `feature-kafka-producer-runtime-contract`(#063)·`feature-idempotency-ownership-protocol-contract`(#070) 의 `/branch-spec` 완료 후 D11·D12 재검토 — 등급: `planned`
|
||
- [ ] **`/depth` 재실행 (최우선)** — 루프 천장에서 종료했고, 마지막 depth Blocking(정지 시점 in-flight 계약)의 처방을 감사 **이후**에 적용해 재검증되지 않았다(§Audit `GATE_CEILING`) — 등급: `needs-confirmation`
|
||
- [ ] stop 계약의 phase 배치·예산을 [[raw/branch-notes/feature-runtime-health-lifecycle-contract]] D4 · [[raw/branch-notes/feature-background-job-async-contract]] owner 와 협의 — 등급: `planned`
|
||
- [ ] inbound leaf 등록·manual ack·rebalance·DLT·inbox 멱등 계약 test 가 통과한다 — 등급: `planned`
|
||
|
||
## 진행 중 메모
|
||
|
||
- 2026-07-28 `/branch-spec` 실행: ca-tmpl ground truth 감사(consumer·inbox 인프라 전부 부재, outbound messaging 은 SPI seam 으로 `actually-implemented`), 자동조사 6건(ack / rebalance / backpressure / poison / retry·DLT / inbox), raw 11건 수집(official-doc 8 + company-tech-blog 3), D1~D14 작성.
|
||
- 2026-07-28 게이트 loop 1: depth `Not ready`(Blocking 2 — D4 선택 축 성립 불가 / 프레임워크 층 ack 공백) + coverage `Not-covered`(Blocking 1 — allowlist 무주인, 위임 주장이 거짓). 조치: D4 재작성, raw 2건 추가 수집(`spring-kafka-ack-mode-...`, `spring-kafka-pause-resume-partitions-...`), **D15 신설**, 위임 정정. → coverage **Covered**(Blocking 0) 달성.
|
||
- 2026-07-28 게이트 loop 2: depth 재감사에서 신규 Blocking 1건(**poll 배치 흡수 규칙 부재** — drop 금지·블로킹 금지·pause 지연이 동시 성립해 합법 행동이 없어지는 구멍) + Should-fix 5건. 조치: 흡수 불변식·ack 발화 지점·모델별 poll 예산 관계식·pause 단위(파티션)·`concurrency` 행·구독 토픽 출처를 각각 명시. 미해소 항목은 §Audit & Findings 의 `SOURCE_GAP-1`·`SOURCE_GAP-5`(잔여).
|
||
|
||
## 결정 사항
|
||
|
||
- 2026-07-28: (D1) 신규 inbound leaf 를 module registry 에 등록해 만든다(19→20). 검토한 대안: 기존 `adapter:outbound:messaging` 에 consumer 를 얹기 — outbound leaf 에 inbound role 이 섞여 기각. / 근거: ca-tmpl `.harness/project/modules.yaml`·`src/settings.gradle` 실측
|
||
- 2026-07-28: (D2) seam 은 유지하되 스켈레톤이 `spring-kafka` 기반 기본 구현을 제공하고, broker 미선택 기동에서는 Kafka auto-config 가 켜지지 않아야 한다. **초안은 현행 코드(`KafkaSender` javadoc "The skeleton carries no Kafka SDK dependency")만 보고 순수 SPI 로 썼다가, 같은 날 작성된 sibling #063 D2 와 축이 갈리는 것을 발견해 재작성했다**(§Audit `SIBLING_DRIFT`). 검토한 대안: 순수 seam-only 유지 — SDK 가 다른 broker 구현을 오염시킨다고 확인되면 그때 후퇴하되 producer 와 함께 결정. / 근거: [[raw/branch-notes/feature-kafka-producer-runtime-contract]] D2
|
||
- 2026-07-28: (D3) `enable.auto.commit=false` + use case 성공·DB 커밋 이후에만 ack. 검토한 대안: 자동 커밋 — 전제조건이 bounded queue 비동기 처리와 충돌해 기각. / 근거: [[raw/official-docs/kafka-consumer-offset-commit-semantics-apache-javadoc]], [[raw/official-docs/kafka-consumer-configs-max-poll-and-commit-defaults]]
|
||
- 2026-07-28: (D4) 순서 단위는 파티션이며 동일 파티션 레코드는 직렬 처리한다(공유 단일 큐 금지). / 근거: [[raw/company-tech-blogs/kafka-multi-tier-retry-topic-dlq-uber]] + outbox D6 + producer [[raw/branch-notes/feature-kafka-producer-runtime-contract]] D5 (`key = aggregateId` 매핑이 `actually-implemented`). 공식 근거 보강은 `SOURCE_GAP-1`
|
||
- 2026-07-28: (D5) backpressure 는 pause/resume 으로 표현하고 레코드를 거부하지 않는다. 검토한 대안: background-job 의 `AbortPolicy` 재사용 — 거부는 유실이라 at-least-once 위반으로 기각. / 근거: [[raw/official-docs/spring-kafka-listener-container-pause-resume-backpressure]]
|
||
- 2026-07-28: (D6) `CooperativeStickyAssignor` pin + `max.poll.*` 명시 pin + `onPartitionsRevoked` 를 유일 커밋 체크포인트로 신뢰 금지. / 근거: [[raw/official-docs/kafka-incremental-cooperative-rebalance-kip429]], [[raw/company-tech-blogs/kafka-consumer-rebalance-cooperative-sticky-verygoodsecurity]]
|
||
- 2026-07-28: (D7) 역직렬화 실패는 리스너 이전 단계에서 감지하고 non-retryable 로 분류한다. 검토한 대안: skip-and-log — 감사 흔적 없이 유실되어 기각. / 근거: [[raw/official-docs/spring-kafka-error-handling-deserializer-poison-record]], [[raw/company-tech-blogs/kafka-poison-pill-consumer-stuck-offset-confluent]]
|
||
- 2026-07-28: (D8) 재시도는 blocking bounded retry 기본, non-blocking retry topic 기각. / 근거: [[raw/official-docs/spring-kafka-non-blocking-retry-topic-ordering-loss]] — "By using this strategy you lose Kafka's ordering guarantees for that topic."
|
||
- 2026-07-28: (D9) dead-letter 발행은 application-core port 경유 — inbound leaf 가 producer 를 직접 보유하지 않는다. / 근거: [[raw/official-docs/spring-kafka-default-error-handler-dlt-fatal-exceptions]] + modules.yaml 의존 제한
|
||
- 2026-07-28: (D10·D11) inbox 기록과 비즈니스 write 를 동일 트랜잭션에서 커밋하고, dedupe key 는 `idempotencyKey` + `eventType` 복합. `(topic,partition,offset)` 단독은 outbox 재발행에서 깨져 기각. / 근거: [[raw/official-docs/idempotent-consumer-microservices-io]], [[raw/branch-notes/feature-domain-event-outbox-contract]] D12·D14
|
||
- 2026-07-28: (D12) owner token(#070) 재사용은 worker fan-out 도입 시로 조건화 — 위임 대상이 스캐폴딩이라 `needs-approval`
|
||
- 2026-07-28: (D13) inbox TTL 수치는 정하지 않고 관계식만 고정 — 외부 근거 부재(`UNSUPPORTED_DECISION`)
|
||
- 2026-07-28: (D14) consumer/inbox 계약 값은 registry 에 없으므로 전부 "신규 제안" 으로만 표기한다. / 근거: ca-tmpl `docs/registries/*.yaml` 실측
|
||
- 2026-07-28: (D15, loop 1 추가) 등록된 `(topic, eventType)` 조합만 소비하고 미등록 조합은 non-retryable 로 회수한다. allowlist 는 코드 handler 등록부로 둔다. 검토한 대안: yaml/env 런타임 등록부(코드와 갈라짐), allowlist 없이 skip(조용한 유실), Schema Registry 위임(스키마 진화만 보고 라우팅을 보지 않음) — 모두 기각. **이 관심사는 초안에서 `feature-schema-serialization-contract` 로 위임한다고 잘못 적었다가 coverage 게이트에서 `missing` 판정을 받아 회수했다**(§Audit `FALSE_DELEGATION`). / 근거: [[raw/official-docs/spring-kafka-default-error-handler-dlt-fatal-exceptions]] `SPRK-ERRH-C1` + outbox D12
|
||
|
||
<!-- section-id: decision-evidence -->
|
||
## Decision Evidence Map / 결정-근거 매핑
|
||
|
||
> 2026-07-28 `/branch-spec` 작성. company-tech-blog 는 `company-case-study` 라벨(공식 best practice 단정 금지). ca-tmpl 코드/registry 대조 결과는 `internal-contract-registry` / `actually-implemented`.
|
||
|
||
| Decision ID | Decision | 선택 조건 (언제 이 결정 / 언제 대안) | Supporting Claims | Evidence Strength | Open Risk |
|
||
|---|---|---|---|---|---|
|
||
| D1 | 신규 inbound leaf `adapter:inbound:messaging-kafka` 를 `.harness/project/modules.yaml` 에 등록(현재 19 → 20). `role: inbound`, package root `dev.caskeleton.adapter.inbound.messaging.kafka`, `allowed_dependencies = [application-core, domain-core, shared-contract]` | 새 inbound transport 는 항상 별도 leaf. **대안(기각)**: 기존 `adapter:outbound:messaging` 에 consumer 를 얹기 — outbound leaf 에 inbound role 이 섞이고 `settings.gradle` 이 registry 를 읽어 include 하므로 role/의존 계약이 흐려진다. 도메인 핸들러가 단 하나뿐인 초소형 fork 라면 leaf 없이 `app-bootstrap` 배선만으로 시작할 수는 있으나, 그때도 registry row 없이 새 소스 디렉터리를 만들 수 없다 | ca-tmpl `.harness/project/modules.yaml` (2026-07-28 실측: modules 19개, inbound 4종 `web`/`grpc`/`graphql`/`websocket` 전수가 동일 3-의존 목록), `src/settings.gradle` (registry 를 읽어 `include` — registry 미등록 모듈은 빌드에 존재조차 못 함) | `internal-contract-registry` + `actually-implemented` | leaf 등록은 `verifyCleanArchitectureDependencies` 와 `CleanArchitectureTest` 의 검사 대상이 늘어나는 것이므로, `mutation_import` 값 선정 등 registry row 의 나머지 필드는 구현 시 기존 inbound row 를 복제해 결정해야 함(§구현 가이드 1 의 `UNSUPPORTED_IMPL_DECISION`) |
|
||
| D2 | consumer 도 producer 와 **같은 축**을 따른다: `KafkaSender` 류 seam 은 유지하되 스켈레톤이 `spring-kafka` 기반 **기본 구현을 제공**하고, `app.messaging.broker` 가 kafka 로 선택되지 않은 기동에서는 Kafka 관련 auto-configuration 이 활성화되지 않아야 한다 | sibling [[raw/branch-notes/feature-kafka-producer-runtime-contract]] D2 가 producer 측에서 이 축을 이미 결정했으므로 consumer 만 순수 SPI 로 남기면 같은 저장소 안에서 두 축이 갈린다. **대안**: Kafka SDK 를 classpath 에 올리는 것이 다른 broker 구현까지 오염시킨다고 확인되면 그때 seam-only 로 후퇴 — 그 판단은 producer 측과 **함께** 내려야 한다(두 결정이 같은 classpath 를 공유) | sibling [[raw/branch-notes/feature-kafka-producer-runtime-contract]] D2 verbatim: "`KafkaSender` seam 은 유지하되 skeleton 이 spring-kafka 기반 기본 구현을 제공하며, **broker 가 선택되지 않은 기동에서는 Kafka 관련 auto-configuration 이 활성화되지 않아야 한다**" + 같은 D2 가 "inbound(consumer)는 신규 leaf 가 필요해 #064 가 모듈 registry migration(19→20)을 소유한다" 로 본 branch 의 D1 을 명시 승인. 현재 코드: `adapter/outbound/messaging/kafka/KafkaSender.java` javadoc "The skeleton carries no Kafka SDK dependency — it is added by the project that enables Kafka." (`actually-implemented` — **아직 SDK-free 상태**), `KafkaAdapterConfig` `@ConditionalOnProperty(app.messaging.broker, havingValue="kafka")` | `internal-cross-reference` (#063 D2 — 그 branch 도 `proposed` 상태) + `actually-implemented` (현행 seam 코드) | **현행 코드(SDK-free)와 #063 D2(기본 구현 제공)가 아직 어긋나 있다** — 본 D2 는 코드가 아니라 sibling 결정에 정렬한 것이고, #063 D2 자체가 `proposed` 다. 또한 #063 D2 의 Open Risk 가 그대로 본 branch 에도 적용된다: classpath 에 SDK 를 올리면 [[raw/branch-notes/feature-capability-provider-selection-contract]] D13("비활성 capability 는 연결·워커·스키마·health contributor 를 만들지 않는다")과 충돌할 수 있고, ca-tmpl 에 `spring.autoconfigure.exclude` 전례가 0건이다. consumer 는 auto-config 가 켜지는 순간 **listener container 가 실제로 broker 에 연결을 시도**하므로 producer 보다 이 충돌이 더 즉각적이다 |
|
||
| D3 | 오프셋 ack 계약: `enable.auto.commit=false`. application use case 성공 **그리고** (비즈니스 write + inbox insert) DB 트랜잭션 커밋이 끝난 뒤에만 오프셋을 커밋한다 | 부작용(비즈니스 write)이 있는 모든 핸들러에서 항상. **대안(기각)**: `enable.auto.commit=true` — 자동 커밋이 at-least-once 를 주기 위한 전제("매 poll 반환분을 다음 poll 전에 전부 소비")가 본 branch 범위의 bounded queue 비동기 처리와 정면 충돌하며, 위반 시 committed offset 이 consumed position 을 앞질러 레코드가 유실된다. 순수 조회(부작용 없음) 핸들러만 있는 토픽이면 이 결정의 위험이 사라지지만 본 branch 범위 밖 | `raw/official-docs/kafka-consumer-offset-commit-semantics-apache-javadoc.md#KAFKA-OFFSET-C2` (자동 커밋의 전제조건 + 위반 시 missing records), `#KAFKA-OFFSET-C3` ("a message should not be considered as consumed until it is completed processing"), `#KAFKA-OFFSET-C4` (insert 후 commit 전 실패 → 재소비 = at-least-once 의 구조), `#KAFKA-OFFSET-C5` (`commitSync` 블로킹 vs `commitAsync` 비블로킹), `raw/official-docs/kafka-consumer-configs-max-poll-and-commit-defaults.md#KAFKA-CONSCFG-C3` (`enable.auto.commit` **기본값 true** — 명시적으로 꺼야 함), `#KAFKA-CONSCFG-C4` (`auto.commit.interval.ms` 기본 5000ms), `raw/official-docs/spring-kafka-ack-mode-manual-commit-and-concurrency.md#SPRK-ACKMODE-C4` ("The default AckMode is BATCH." — 프레임워크 층에서도 **명시 설정하지 않으면 자동 배치 커밋**), `#SPRK-ACKMODE-C2` (`MANUAL_IMMEDIATE` = acknowledge() 호출 즉시 커밋), `#SPRK-ACKMODE-C1` (`MANUAL` 은 이후 `BATCH` 시맨틱), `#SPRK-ACKMODE-C3` (리스너가 `AcknowledgingMessageListener` 여야 함) | `official-vendor-doc` | `commitSync` vs `commitAsync` 선택은 미확정 — 동기 커밋은 처리량을 깎고 비동기 커밋은 실패가 콜백으로만 전달된다(`KAFKA-OFFSET-C5`). §구현 가이드 2 의 `UNSUPPORTED_IMPL_DECISION` |
|
||
| D4 | 소비 병렬성의 **순서 단위는 파티션**이다. 동일 파티션의 레코드는 항상 직렬로 처리하며, 여러 파티션의 레코드를 하나의 공유 큐/워커풀에 섞는 구현은 금지. **스켈레톤 기본 형태는 "poll 스레드 + 파티션별 독립 bounded queue + 파티션당 직렬 워커"** 다 | 선택 축은 동시성 크기가 **아니다** — 직렬 제약 때문에 어느 형태든 병렬도 상한은 파티션 수로 같다. 진짜 축은 **poll 스레드를 처리 지연에서 분리할 필요가 있는가**다. 핸들러 처리시간이 `max.poll.interval.ms` 예산 안에서 끝난다고 보장할 수 없거나(외부 I/O 포함) 재시도 backoff 가 그 예산을 잠식하면 → 큐 분리(기본). **대안**: 핸들러가 순수 CPU·단일 DB write 로 짧고 p99 가 예측 가능하면 리스너 인라인 처리로 단순화 가능 — 이때는 큐가 없으므로 D5 의 pause 도 불필요해지고, 대신 `max.poll.records` 를 낮춰 poll 예산을 지킨다. 순서 무관 이벤트만 싣는 토픽이라도 공유 큐는 기본이 아니며 토픽별 명시 선언이 필요하다 | `raw/company-tech-blogs/kafka-multi-tier-retry-topic-dlq-uber.md#UBER-REPROC-C2` ("Kafka only guarantees in-order processing within partitions and not across them" — 순서 단위가 파티션이라는 사실), [[raw/branch-notes/feature-domain-event-outbox-contract]] D6 (per-aggregate FIFO 를 보장 단위로 정의, global ordering 미보장), [[raw/branch-notes/feature-kafka-producer-runtime-contract]] D5 (순서를 파티션 단위로만 주장하고 `key = aggregateId` 로 per-aggregate FIFO 에 대응 — `OutboxMessagePublishAdapter` 의 `topic=eventType, key=aggregateId` 가 `actually-implemented` 로 확인됨, 2026-07-28), `raw/official-docs/spring-kafka-ack-mode-manual-commit-and-concurrency.md#SPRK-ACKMODE-C7` ("If the concurrency is greater than the number of TopicPartitions, the concurrency is adjusted down such that each container gets one partition." — **프레임워크가 병렬도를 파티션 수로 잘라낸다**는 공식 근거로, "동시성을 키워도 파티션 수가 상한" 이라는 D4 의 선택 축 재정의를 뒷받침) | `company-case-study` (UBER-REPROC-C2) + `internal-cross-reference` (outbox D6 · producer D5) + `actually-implemented` (key 매핑) | **"파티션 하나는 그룹 내 정확히 한 consumer 가 소비한다"는 Kafka 공식 verbatim 을 아직 수집하지 못했다** — 현재 이 사실의 직접 근거는 회사 블로그 1건뿐이다(§Audit `SOURCE_GAP-1`). 또한 producer D5 의 Open Risk 가 그대로 전이된다: **파티션 수를 늘리면 같은 `aggregateId` 가 다른 파티션으로 가서 per-aggregate FIFO 가 깨지고**, 그 순간 본 D4 의 직렬 처리 단위도 의미를 잃는다 |
|
||
| D5 | backpressure 는 bounded queue 포화 시 **해당 파티션만** `pausePartition()`, drain 후 `resumePartition()` 로 처리한다. 레코드를 거부(drop)하거나 예외로 버리지 않으며, poll 스레드를 블로킹하지도 않는다. 이 셋이 동시에 성립하려면 **pause 요청 시점에 직전 poll 배치를 흡수할 큐 여유가 남아 있어야** 한다(§구현 가이드 2 의 흡수 불변식) | 항상. **대안(기각)**: sibling [[raw/branch-notes/feature-background-job-async-contract]] D7·D8 의 `AbortPolicy`(거부 후 `JOB_EXECUTOR_REJECTED` 로그)를 그대로 적용 — 거부는 레코드 유실이므로 at-least-once 계약(상속 결정)을 깬다. 다운스트림 장기 장애로 pause 로도 흡수가 안 되면 D8 재시도 → D9 회수 경로 | `raw/official-docs/spring-kafka-listener-container-pause-resume-backpressure.md#SPRK-PAUSE-C2` ("When a container is paused, it continues to poll() the consumer, avoiding a rebalance if group management is being used, but it does not retrieve any records" — pause 가 그룹 이탈을 유발하지 않는 근거), `#SPRK-PAUSE-C1` (pause 는 다음 poll 직전, resume 은 현재 poll 반환 직후 반영), `#SPRK-PAUSE-C3` (`pauseImmediate` 기본 false = 이전 poll 분 처리 완료 후 적용), `#SPRK-PAUSE-C4` (`isConsumerPaused()` 라야 실제 정지 확인), `raw/official-docs/spring-kafka-pause-resume-partitions-on-listener-containers.md#SPRK-PAUSEPART-C1` (2.7~ `pausePartition(TopicPartition)`/`resumePartition(TopicPartition)` — **파티션 단위** pause 가 D4 의 파티션별 큐와 짝을 이룬다), `#SPRK-PAUSEPART-C2` (poll 경계 반영), `#SPRK-PAUSEPART-C3` (요청 vs 실제 정지 구분), `#SPRK-PAUSEPART-C4` (`ConsumerPartitionPausedEvent`/`ResumedEvent` — 관측 지점) | `official-vendor-doc` | **rebalance·재배정 시 파티션 pause 상태가 유지되는지 초기화되는지는 어느 문서도 말하지 않는다**(`#SPRK-PAUSEPART-C5` 부재 확인 — 파티션 pause 페이지·컨테이너 pause 페이지 모두 rebalance 어휘 자체가 없다) → §구현 가이드 3 의 `onPartitionsAssigned` 행을 `UNSUPPORTED_IMPL_DECISION` 으로 라벨링했다. 또한 pause 의 안전성은 **poll 루프 스레드가 블로킹되지 않는다**는 전제 위에서만 성립한다 |
|
||
| D6 | rebalance 안정성: `partition.assignment.strategy` 를 `CooperativeStickyAssignor` 로 pin 하고, `max.poll.interval.ms`/`max.poll.records` 를 실측 처리시간 기준으로 명시 pin 한다. `onPartitionsRevoked` 를 rebalance 시작 신호나 **유일한** 커밋 체크포인트로 신뢰하지 않는다 | 파티션·컨슈머 수가 많고 배포로 rebalance 가 잦은 배포에서 cooperative. **대안**: 파티션이 한 자릿수이고 그룹이 안정적이면 기본값(`[RangeAssignor, CooperativeStickyAssignor]`) 을 그대로 두어도 되며, 이 기본 목록 덕분에 나중에 `RangeAssignor` 만 제거하는 **단일 rolling bounce** 로 전환할 수 있다 | `raw/official-docs/kafka-incremental-cooperative-rebalance-kip429.md#KIP429-C2` (EAGER = rebalance 전 소유 파티션 전량 revoke), `#KIP429-C3` (COOPERATIVE = 소유 파티션 유지), `#KIP429-C5` ("it is possible for #onPartitionsRevoked to never be invoked at all during a rebalance, and should not be relied on to signal that a rebalance has started"), `#KIP429-C1` (Accepted, Kafka 2.4.0), `#KIP429-C4` (`onPartitionsLost` 의미), `raw/official-docs/kafka-consumer-configs-max-poll-and-commit-defaults.md#KAFKA-CONSCFG-C1` (`max.poll.interval.ms` 기본 300000ms, 초과 시 실패 간주 + rebalance), `#KAFKA-CONSCFG-C2` (`max.poll.records` 기본 500), `#KAFKA-CONSCFG-C7` (기본 전략 목록 + 단일 rolling bounce 업그레이드), `raw/company-tech-blogs/kafka-consumer-rebalance-cooperative-sticky-verygoodsecurity.md#VGS-REBAL-C2`·`#VGS-REBAL-C3`·`#VGS-REBAL-C4` (운영 사례) | `official-vendor-doc` (KIP-429 + consumer configs) + `company-case-study` (VGS — 100 consumer/aiokafka 규모 전제가 다름, best practice 로 격상 금지) | 구체 pin 값(`max.poll.interval.ms` 를 얼마로) 은 처리시간 실측 없이 정할 수 없다 — §구현 가이드 3 의 `UNSUPPORTED_IMPL_DECISION`. D5 의 pause 기반 backpressure 를 쓰면 poll 이 계속되므로 이 값의 압박은 줄지만, poll 스레드가 블로킹되는 순간 동일 실패로 되돌아간다 |
|
||
| D7 | poison/역직렬화 실패는 **리스너 호출 이전 단계(deserializer 경계)** 에서 감지하고 `non-retryable` 로 분류해 첫 실패에 곧바로 회수 경로(D9)로 보낸다 | 구조적 파싱 실패(스키마 불일치·깨진 바이트)일 때. **대안/보완**: 파싱은 성공했지만 도메인 규칙(허용 handler·schema·version allowlist)을 위반하는 "의미상 poison" 은 이 경로가 아니라 application 경계의 예외 분류로 다룬다 — 두 실패는 발생 위치가 달라 상호 배타가 아니라 역할 분담이다 | `raw/official-docs/spring-kafka-error-handling-deserializer-poison-record.md#SPRK-EHD-C1` (역직렬화 실패는 `poll()` 반환 이전에 발생해 리스너 레벨에서 처리 불가), `#SPRK-EHD-C2` (실패 시 null + `DeserializationException` 헤더 with cause + raw bytes), `#SPRK-EHD-C3` (헤더가 있으면 컨테이너의 ErrorHandler 호출, "The record is not passed to the listener"), `raw/official-docs/spring-kafka-default-error-handler-dlt-fatal-exceptions.md#SPRK-ERRH-C1` (`DeserializationException` 등 6종을 기본 fatal 로 분류 — "since these exceptions are unlikely to be resolved on a retried delivery"), `raw/company-tech-blogs/kafka-poison-pill-consumer-stuck-offset-confluent.md#CONF-POISON-C3`·`#CONF-POISON-C4` (미처리 시 offset 정체 + 무한 고속 재시도) | `official-vendor-doc` (Spring reference 2종) + `company-case-study` (Confluent 벤더 블로그 — 메커니즘 설명, 타사 운영 사례 아님) | 네트워크 truncation 처럼 **실제로는 일시적인데 역직렬화 실패로 나타나는** 엣지가 fatal 로 오분류된다(sibling [[raw/branch-notes/feature-outbound-http-client-baseline]] D12 의 "4xx 일괄 PERMANENT 분류" 와 동형 미해결). 이 엣지의 처리는 구현 시 결정 |
|
||
| D8 | 재시도 기본값은 **blocking bounded retry** — 같은 파티션에서 backoff 재시도하고 순서를 보존한다. non-blocking retry topic 체인(`topic-retry-N`)은 기본 기각. backoff/max attempts 어휘는 sibling 에 위임 | 순서 보장(D4)이 요구되는 토픽이면 blocking. **대안**: 특정 토픽이 순서 무관 이벤트만 싣고 처리량이 최우선이면 그 토픽에 한해 retry topic 채택 — 단 "이 토픽은 순서를 포기한다" 를 명시 선언해야 한다. 또한 backoff 총합이 `max.poll.interval.ms` 를 넘길 위험이 있으면 스레드 정지형이 아니라 **컨테이너 pause 형 backoff** 를 쓴다 | `raw/official-docs/spring-kafka-non-blocking-retry-topic-ordering-loss.md#SPRK-RETRYTOPIC-C3` ("By using this strategy you lose Kafka's ordering guarantees for that topic." — 공식 자인), `#SPRK-RETRYTOPIC-C1` (retry topic 은 back-off timestamp 로 파티션 소비를 일시 중지), `#SPRK-RETRYTOPIC-C2` (소진 시 DLT), `raw/official-docs/spring-kafka-default-error-handler-dlt-fatal-exceptions.md#SPRK-ERRH-C2` (기본 backoff 는 consumer 스레드를 정지시키며, 지연이 `max.poll.interval.ms` 보다 길 때를 위해 `ContainerPausingBackOffHandler` 제공), `raw/company-tech-blogs/kafka-multi-tier-retry-topic-dlq-uber.md#UBER-REPROC-C1`·`#UBER-REPROC-C2`·`#UBER-REPROC-C4` (다단계 retry topic 사례 — 순서 비보장 수용이 전제), [[raw/branch-notes/feature-background-job-async-contract]] D4 (exponential backoff with jitter / max attempts 3 / DLQ after exhausted — 어휘 위임) | `official-vendor-doc` (Spring reference 2종) + `company-case-study` (Uber — 순서 비보장 전제가 본 계약과 다름) + `internal-cross-reference` (backoff 어휘) | background-job D4 의 "DLQ after exhausted" 는 그 branch 에서 아직 외부 근거가 없는 항목이다(위임 대상의 잔여 `UNSUPPORTED`). 또한 blocking retry 는 실패가 잦아지면 해당 파티션 전체를 정체시킨다 — 실패율 임계와 pause 전환 기준은 미확정(§구현 가이드 4) |
|
||
| D9 | 최종 실패(재시도 소진 또는 D7 non-retryable)는 dead-letter 로 회수하고 감사 흔적을 남긴다. **발행은 inbound leaf 가 직접 producer 를 들지 않고 `application-core` 의 outbound port 를 경유**한다 | 항상. **대안**: dead-letter 를 Kafka 토픽이 아니라 **DB row 로만** 표현하면 producer 자체가 불필요해 모듈 경계 문제가 사라진다(outbox `SKIP LOCKED` 선례 재사용). 외부 시스템이 DLT 토픽을 직접 구독해야 하는 요구가 있으면 토픽 방식, 내부 운영자만 조회하면 DB row 방식 | `raw/official-docs/spring-kafka-default-error-handler-dlt-fatal-exceptions.md#SPRK-ERRH-C4` ("The recoverer requires a KafkaTemplate<Object, Object>, which is used to send the record." — DLT 발행에 producer 필수), `#SPRK-ERRH-C3` (기본 명명 `<originalTopic>-dlt` + 동일 partition + partition 수 요건), `#SPRK-ERRH-C5` (recoverer 미구성 시 기본은 로그만 — DLT 로 안 감), ca-tmpl `.harness/project/modules.yaml` (inbound leaf 의 `allowed_dependencies` 에 outbound leaf 가 **없음** — 2026-07-28 실측), `raw/company-tech-blogs/kafka-multi-tier-retry-topic-dlq-uber.md#UBER-REPROC-C5` (DLQ = 지속 실패의 종착점) | `official-vendor-doc` + `internal-contract-registry` (모듈 경계) + `company-case-study` | **감사된 replay(누가·언제·무엇을 재처리했는지)를 규정하는 외부 근거는 어디에도 없다** — Spring/Confluent/Uber 어느 문서도 replay audit trail 을 다루지 않는다. audit 기록 스키마는 `UNSUPPORTED_IMPL_DECISION`(§구현 가이드 4) |
|
||
| D10 | `InboxStorePort`(application-core, **신규 — 현재 코드 부재**) 에 처리한 메시지 식별자를 기록하고, 그 insert 를 비즈니스 write 와 **동일 `TransactionPort.inWrite` 경계**에서 커밋한다. ack 는 그 커밋 성공 이후(D3) | 이벤트가 여러 aggregate/부작용에 걸치거나 감사·replay 가시성이 필요할 때(본 branch 범위가 "감사된 replay" 를 포함하므로 기본). **대안**: 이벤트가 정확히 하나의 aggregate row 를 1회성으로 만들고 그 row 에 자연 유니크 키가 있으면 별도 inbox 없이 비즈니스 엔티티 자체에 ID 를 저장하는 변형으로 대체 가능 — 대신 처리 이력 조회를 포기 | `raw/official-docs/idempotent-consumer-microservices-io.md#MSIO-IDEMPC-C2` (처리한 메시지 ID 를 DB 에 기록해 멱등), `#MSIO-IDEMPC-C3` ("After starting the database transaction, the message handler inserts the message's ID into the PROCESSED_MESSAGE table."), `#MSIO-IDEMPC-C4` (복합 PK 위반으로 duplicate INSERT 실패), `#MSIO-IDEMPC-C5` (비즈니스 엔티티 저장 변형), `raw/official-docs/microservices-io-transactional-outbox.md#MSIO-OUTBOX-C5` ("The solution is for the service that sends the message to first store the message in the database as part of the transaction that updates the business entities." — producer 측 동일-트랜잭션 원리의 명시 앵커. consumer 측 대칭 적용의 근거로 인용하되 원문은 producer 문맥임을 유지), [[raw/branch-notes/feature-domain-event-outbox-contract]] D7 (consumer 는 at-least-once + idempotencyKey dedupe 의무 — 메커니즘은 consumer branch 위임), ca-tmpl `application/transaction/TransactionPort.java` (`actually-implemented`) | `engineering-blog` (microservices.io = Chris Richardson 개인 패턴 카탈로그 — 벤더 공식 아님) + `internal-cross-reference` + `actually-implemented` (TransactionPort) | **"메시지 ID 기록과 비즈니스 데이터 갱신이 같은 트랜잭션이어야 한다" 는 명시 문장은 consumer 측 인용 원문에 없다**(원 raw 의 Usage Boundaries 에 기록됨) — 원문은 handler 가 트랜잭션을 시작해 ID 를 INSERT 한다는 것까지만 말한다. 동일 트랜잭션 요구는 `MSIO-OUTBOX-C5`(producer 측 동일-트랜잭션 원리)의 대칭 적용이라는 **해석**이다. 그 해석이 필요한 이유는 연역으로 닫힌다: **inbox insert 만 커밋되고 비즈니스 write 가 롤백되면 그 메시지는 이후 영구히 "처리됨" 으로 스킵된다**(조용한 누락). 반대로 비즈니스 write 만 커밋되면 중복 실행이 남는다 — 두 실패 모두 트랜잭션을 합쳐야만 사라진다 |
|
||
| D11 | dedupe key 는 producer envelope 의 `idempotencyKey`(= `eventId`, ULID) 를 1차로 쓰고, inbox 유니크 제약은 `(idempotencyKey, eventType)` 복합으로 건다. `(topic, partition, offset)` 은 **유일 dedupe key 로 채택하지 않는다** | 항상. **대안 기각 이유**: outbox relay 가 publish 성공 후 status 갱신 전 crash 하면 **같은 논리 이벤트가 다른 offset 으로 재발행**되므로 `(topic,partition,offset)` 는 그것을 서로 다른 이벤트로 오판한다. 반대로 토픽 재생성/DR 미러링 시엔 offset 이 0부터 재할당되어 정상 이벤트를 중복으로 오판할 수 있다. 복합 키를 쓰는 이유는 CloudEvents 가 dedup 단위를 `source + id` **조합**으로 규정하는 것과 같은 취지 | [[raw/branch-notes/feature-domain-event-outbox-contract]] D12 (envelope required fields = `eventId`/`occurredAt`/`aggregateId`/`eventType`/`correlationId`/`idempotencyKey`, 구현상 `idempotencyKey = eventId` ULID), 동 D14 (outbox idempotencyKey 는 API `Idempotency-Key` 와 별개 scope), 동 §엣지 ("publish 성공 후 status 갱신 전 crash → 동일 event 재발행"), `raw/official-docs/cloudevents-spec-required-attributes.md#CLOUDEVT-C2` (`source`+`id` 조합이 고유성 단위, consumer 는 동일 조합을 duplicate 로 간주 가능), `raw/official-docs/idempotent-consumer-microservices-io.md#MSIO-IDEMPC-C4` (복합 PK 로 중복 탐지) | `internal-cross-reference` (outbox 계약) + `official-vendor-doc` (CloudEvents) + `engineering-blog` (MSIO) | 이 key 는 **producer 가 동일 논리 이벤트에 항상 같은 `eventId` 를 재사용한다**는 전제에 의존한다. outbox branch 에 그 보장을 명시한 D-row 는 없다(2026-07-28 확인) — 어긋나면 정상 이벤트가 조용히 누락된다. `(idempotencyKey, eventType)` 복합 제약이 그 오류를 탐지하는 최소 방어선 |
|
||
| D12 | owner token 프로토콜(claim/renew/complete/release, #070) 재사용은 **조건부**다 — 파티션당 직렬 처리(D4 기본)에서는 insert-once inbox 로 충분하고, worker fan-out 을 도입해 rebalance 중 zombie consumer 가 같은 메시지를 동시 처리할 수 있게 되면 그때 claim/lease 를 재사용한다 | 동기 직렬 처리 → 단순 inbox. worker fan-out + rebalance 노출 → owner token(`SAME_STORE_TRANSACTIONAL` 등급). 판단은 D4 의 concurrency 모델 확정 이후 | [[raw/branch-notes/feature-idempotency-ownership-protocol-contract]] §범위 (claim 결과 `ACQUIRED`/`REPLAY`/`IN_PROGRESS`/`FINGERPRINT_MISMATCH`, 실행 lease 와 replay TTL 분리, 보증 등급 3종 — **D-row 는 아직 미작성**, 2026-07-28 확인), `raw/official-docs/kafka-incremental-cooperative-rebalance-kip429.md#KIP429-C4` (`onPartitionsLost` = 이미 소유권을 잃은 상태의 콜백 → zombie 구간 존재의 근거) | `internal-cross-reference` (위임 대상이 스캐폴딩 상태 — `needs-confirmation`) + `official-vendor-doc` (KIP429-C4) | 위임 대상 #070 이 스캐폴딩(D-row 0개)이라 재사용 비용을 아직 평가할 수 없다. #070 의 `/branch-spec` 완료 후 본 D-row 재검토 필요 — 그 전에는 `needs-approval` |
|
||
| D13 | inbox row 보존 기간(TTL)의 **구체 수치는 정하지 않는다**. 원칙만 고정: (a) 토픽 `retention.ms` 와 (b) 지원하려는 최대 dead-letter replay 창 중 **긴 쪽 이상**, 그리고 (c) 무한 보관 금지(reaper 필요). 이에 따라 **dedupe 보증은 TTL 창 안의 재전달에 한정**되며, 창 밖 재전달은 신규 처리로 간주된다 — 이 한계를 계약으로 명시한다 | 항상 원칙 적용. 수치는 운영 환경의 retention/replay 창이 확정된 뒤 결정 | **UNSUPPORTED_DECISION** — 외부 근거 부재. 이번 조사에서 CloudEvents spec / AWS Prescriptive Guidance / microservices.io 어디에도 inbox TTL 수치를 규정한 문장이 없음을 확인했다. 실무 관행값(Stripe 계열 ~72시간, Toss 15일)은 **API-level idempotency-key 도메인**의 사례라 consumer inbox 로 직접 이전할 수 없다. 사용자 trade-off: TTL 이 replay 창보다 짧으면 replay 된 메시지가 신규 처리로 오판되므로, 수치를 지어내는 것보다 관계식만 고정하는 편이 안전하다 | `internal-policy` (외부 근거 없음 명시) | 절대 상한이 없으면 "무한 보관 금지" 가 실질적으로 집행되지 않는다 — reaper 주기와 상한은 구현 시 outbox reaper(`ca-skeleton.outbox.published-retention` 선례)를 모델로 결정 |
|
||
| D15 | **handler/schema/version allowlist**: 이 leaf 는 등록된 **`(topic, eventType)`** 조합만 소비하고(`schemaVersion` 축은 envelope 에 그 필드가 생긴 뒤 추가 — Open Risk 참조), 미등록 조합은 처리하지 않고 D7 의 non-retryable 경로로 회수한다. allowlist 는 **코드에 선언된 handler 등록부**(핸들러가 자기 `(topic, eventType, 지원 schemaVersion 범위)` 를 선언하고 기동 시 조합의 중복·공백을 검증)이며, 별도 런타임 설정 파일이나 registry yaml 로 두지 않는다 | 항상 — 미등록 조합을 조용히 무시하거나(유실) 아무 handler 에나 라우팅하는 것(오처리)이 둘 다 금지되므로 명시 allowlist 가 필요하다. **대안(기각)**: (a) env/yaml 기반 런타임 allowlist — 코드의 handler 와 설정이 갈라져 "등록했는데 handler 가 없는" 상태가 런타임에만 드러난다. (b) allowlist 없이 미등록 조합을 skip — 유실이 조용해져 at-least-once 계약의 관측성을 깬다. (c) Schema Registry 의 호환성 검사에 위임 — 그 검사는 *payload 스키마 진화*를 보고 *어느 handler 가 이 이벤트를 맡는가*를 보지 않는다 | **UNSUPPORTED_DECISION (부분)** — allowlist 의 *존재 필요성* 은 근거가 있다: `raw/official-docs/spring-kafka-default-error-handler-dlt-fatal-exceptions.md#SPRK-ERRH-C1` 이 `MethodArgumentResolutionException`·`NoSuchMethodException`·`ClassCastException`(= 라우팅/시그니처 불일치 계열)을 **fatal** 로 분류해 재시도 대상이 아님을 확정하고, `raw/company-tech-blogs/kafka-poison-pill-consumer-stuck-offset-confluent.md#CONF-POISON-C1` 이 "항상 실패하는 레코드" 개념을 정의한다. envelope 의 `eventType` 은 [[raw/branch-notes/feature-domain-event-outbox-contract]] D12 가 required 로 확정. **그러나 "코드 등록부 vs 설정 파일" 이라는 형태 선택과 `schemaVersion` 필드의 존재 자체는 외부 근거가 없다** — envelope required 6필드에 `schemaVersion` 은 **없다**(outbox D12 실측). 사용자 trade-off: 설정과 코드가 갈라지는 실패를 없애려면 등록부를 코드에 두는 편이 안전하고, 버전 축은 필요해질 때 envelope 확장으로 추가한다(지금 발명하지 않음) | `official-vendor-doc` (fatal 분류) + `internal-cross-reference` (envelope) + `internal-policy` (형태 선택 — 근거 없음) | **`schemaVersion` 이 현재 envelope 에 없다** — 이 축을 실제로 쓰려면 outbox D12 의 required 필드 확장이 필요하고 그것은 [[raw/branch-notes/feature-domain-event-outbox-contract]] 소유다. 확장 전까지 allowlist 의 실효 키는 `(topic, eventType)` 2축뿐이다. 또한 이 관심사는 노트 초안에서 [[raw/branch-notes/feature-schema-serialization-contract]] 로 위임했다고 적었으나 **그 branch 는 이 관심사를 소유하지 않음**(2026-07-28 grep 재확인 — handler/topic/routing 언급 0건, Avro/JSON 직렬화 전용). §Audit `FALSE_DELEGATION` 참조 |
|
||
| D14 | consumer/inbox 용 error code·metric·env key 는 ca-tmpl registry 에 **하나도 등록돼 있지 않다**. 본 노트는 전부 "신규 제안" 으로만 표기하고 기존 값처럼 단정하지 않는다. category 는 반드시 `Category.java` 의 10종 안에서 고른다 | 항상. registry 반영은 구현 branch 착수 시 `owner_branch: feature-kafka-consumer-inbox-contract` 로 등록 | ca-tmpl `docs/registries/error-codes.yaml` (2026-07-28 grep: `OUTBOX_*`/`JOB_*` 는 있으나 consumer/inbox row **부재**), `metrics.yaml` (`outbox.*`/`job.*` 만 존재), `env-keys.yaml` (`APP_MESSAGING_BROKER`/`APP_MESSAGING_KAFKA_BROKERS` 만 존재, consumer 키 부재 — owner 는 `feature-integration-adapter-templates`), `src/shared-contract/.../error/Category.java` (10-value enum, 코드 SSOT) | `internal-contract-registry` + `actually-implemented` (Category enum) | 신규 코드/메트릭 명명 자체는 외부 근거가 없다 — §구현 가이드 5 의 `UNSUPPORTED_IMPL_DECISION`. registry 의 `required_test` / `runbook_link` 필드는 retryable=true 행에 runbook 을 강제하므로 제안 시 runbook 작성 의무가 따라온다 |
|
||
|
||
<!-- section-id: implementation -->
|
||
## 구현 가이드
|
||
|
||
> *결정(D-row)* 이 "*무엇*" 이라면 본 §는 "*어디에 어떻게*" 의 사전 명세. 모든 cell 은 Decision ID + Supporting Claim 의 도출(R1). 근거 없는 detail 은 `UNSUPPORTED_IMPL_DECISION`(R2). 본 branch 결정 범위 밖 detail 은 두지 않음(R3).
|
||
>
|
||
> **증거 등급 주의**: 아래 "위치" 열의 클래스·포트는 2026-07-28 ca-tmpl `src/` grep 결과 **전부 부재**다. 기존 코드(`TransactionPort`, `MessageBroker`, `KafkaSender`, `Category`)만 `actually-implemented` 이고 나머지는 `planned` 이다.
|
||
|
||
### 1. 모듈 등록·배치 (D1, D2)
|
||
|
||
> **Trace**: D1 (ca-tmpl `.harness/project/modules.yaml` 실측 — modules 19개, inbound 4종 전수 동일 3-의존; sibling #063 D2 가 "모듈 registry migration(19→20)은 #064 소유" 로 명시 승인) + D2 (#063 D2 와 같은 축 + 현행 `KafkaSender`/`MessageBroker` javadoc).
|
||
>
|
||
> - **UNSUPPORTED_IMPL_DECISION**: (1) registry row 의 `mutation_import` 값 — 기존 inbound 4종은 전부 `dev.caskeleton.adapter.outbound.persistence.OutboxStoreAdapter` 를 쓰지만 그 선정 근거는 문서화돼 있지 않다. trade-off: 선례를 그대로 복제하는 편이 새 값을 발명하는 것보다 안전. (2) leaf slug 를 `messaging-kafka` 로 둘지 `kafka` 로 둘지 — outbound 는 `messaging` 아래 broker 를 두는 2단 구조(`messaging/kafka/`)인데 inbound 는 1단(`web`/`grpc`)이다. trade-off: project §8.0 이 `adapter:inbound:messaging-kafka` 를 이미 명시했으므로 그 이름을 따른다(발명 아님).
|
||
|
||
| 항목 | 위치 (module / path) | 증거 등급 | Trace |
|
||
|---|---|---|---|
|
||
| module registry row 추가 (`id: adapter-inbound-messaging-kafka`, `role: inbound`, `gradle_path: :adapter:inbound:messaging-kafka`, `source_path: src/adapter/inbound/messaging-kafka`, `allowed_dependencies: [application-core, domain-core, shared-contract]`, `claude: src/adapter/inbound/messaging-kafka/CLAUDE.md`) | ca-tmpl `.harness/project/modules.yaml` | `planned` (현재 19 modules) | D1 |
|
||
| Gradle include — **별도 작업 없음**. `src/settings.gradle` 이 registry 를 읽어 `include` 하므로 registry row 추가만으로 모듈이 생긴다 | ca-tmpl `src/settings.gradle` | `actually-implemented` | D1 |
|
||
| leaf 로컬 규칙 문서 (`CLAUDE.md`) — 기존 inbound leaf 와 동일 골격(Registered identity / Responsibility / Boundaries / Tests) | `src/adapter/inbound/messaging-kafka/CLAUDE.md` | `planned` | D1 |
|
||
| consumer seam + 스켈레톤 기본 구현 — outbound `KafkaSender` 대칭 seam 을 두되 `spring-kafka` 기반 기본 listener 구현을 스켈레톤이 제공 | `adapter:inbound:messaging-kafka` `dev/caskeleton/adapter/inbound/messaging/kafka/` | `planned` | D2 |
|
||
| 활성화 게이트 — `app.messaging.broker=kafka` 가 아니면 consumer 구성이 **켜지지 않아야** 한다(연결·listener container 미생성). outbound `KafkaAdapterConfig` 의 `@ConditionalOnProperty(app.messaging.broker, havingValue="kafka")` 패턴을 따른다 | `adapter:inbound:messaging-kafka` config | `planned` (패턴 자체는 outbound 에 `actually-implemented`) | D2 |
|
||
| 대칭 선례 (변경 없음, 참조만) — `MessageBroker` SPI + `KafkaSender` + `MessagingConfig` 중앙 조립 + `app.messaging.broker` 선택 | `adapter:outbound:messaging` `dev/caskeleton/adapter/outbound/messaging/` | `actually-implemented` | D2 |
|
||
|
||
### 2. 소비 파이프라인과 ack 순서 (D3, D5, D10)
|
||
|
||
> **Trace**: D3 (`KAFKA-OFFSET-C2`/`C3`/`C4`, `KAFKA-CONSCFG-C3`/`C4`) + D5 (`SPRK-PAUSE-C1`~`C4`) + D10 (`MSIO-IDEMPC-C3`/`C4`) + ca-tmpl `TransactionPort` (`actually-implemented`).
|
||
>
|
||
> - **UNSUPPORTED_IMPL_DECISION**: (1) `commitSync` vs `commitAsync` — `KAFKA-OFFSET-C5` 는 둘의 시맨틱만 말하고 어느 쪽을 쓰라고 권고하지 않는다. trade-off: 커밋 실패가 조용히 삼켜지면 재처리 폭이 커지므로 **동기 커밋을 기본**으로 두고, 처리량 문제가 실측되면 비동기 + 실패 임계 카운터로 전환. (2) 큐 용량·pause 임계(high/low watermark) 수치 — 근거 없음. trade-off: 값 자체보다 "큐 포화가 pause 로 이어진다"는 관계를 계약으로 고정하고 수치는 설정으로 노출. (3) **D2 가 고른 프레임워크 층(spring-kafka)에서 ack 를 무엇으로 표현하고 어느 스레드에서 호출하는가** — 아래 "ack 메커니즘" 표 참조. D4 가 poll 스레드와 워커를 분리하므로 "워커 스레드에서 ack 를 호출해도 되는가" 가 곧바로 문제가 된다.
|
||
|
||
**ack 메커니즘 (D2 의 프레임워크 층)** — 이 work item 의 완료 조건이 "manual ack 계약 test 통과" 이므로 공백으로 둘 수 없다:
|
||
|
||
| 항목 | 명세 | 근거 |
|
||
|---|---|---|
|
||
| ack 모드 | `AckMode` 를 **명시적으로** `MANUAL_IMMEDIATE` 로 설정한다. **설정을 빠뜨리면 기본값이 `BATCH`** 라 poll 배치 단위 자동 커밋으로 조용히 되돌아간다 — D3 위반이 침묵으로 발생하는 지점 | `SPRK-ACKMODE-C4` ("The default AckMode is BATCH."), `SPRK-ACKMODE-C2` (`MANUAL_IMMEDIATE` = acknowledge() 호출 즉시 커밋) |
|
||
| `MANUAL` 을 쓰지 않는 이유 | `MANUAL` 은 acknowledge() 이후 `BATCH` 와 같은 시맨틱(= poll 반환분 전체 처리 후 커밋)이라 "DB 커밋 직후 그 레코드만 ack" 라는 D3 의 시점 계약을 표현하지 못한다 | `SPRK-ACKMODE-C1` |
|
||
| 리스너 타입 제약 | `MANUAL`/`MANUAL_IMMEDIATE` 는 리스너가 `AcknowledgingMessageListener`(또는 배치형)여야 한다 — 리스너 시그니처가 이 결정에 종속된다 | `SPRK-ACKMODE-C3` |
|
||
| ack 호출 스레드 | **poll/리스너 스레드에서 ack 한다.** 워커에서 직접 ack 하지 않고 완료 offset 을 리스너 스레드로 되돌린다 | `SPRK-ACKMODE-C5` (`nack()` 은 리스너를 호출한 consumer 스레드에서만 호출 가능), `SPRK-ACKMODE-C6` (부분 배치 acknowledge 도 "리스너 스레드에서 호출되어야 한다") — `Acknowledgment` 의 최소 두 메서드가 스레드에 묶여 있으므로 워커 스레드 ack 를 문서 근거 없이 가정하지 않는다 |
|
||
| 커밋 순서 | 파티션별 **완료 연속 구간의 최솟값**까지만 ack — 워커 완료 순서가 poll 순서와 어긋나도 offset 이 앞질러 가지 않는다 | D3(`KAFKA-OFFSET-C2` 의 "committed offset 이 consumed position 을 앞지르면 유실") + 위 스레드 제약 |
|
||
| ack 발화 지점 | 워커는 완료분을 파티션별 pending 구조에 적재하되 **offset 뿐 아니라 ack 수단(레코드의 `Acknowledgment` 핸들)까지 함께 보관**한다 — 커밋은 `Acknowledgment.acknowledge()` 호출로만 일어나므로(`SPRK-ACKMODE-C2`) offset 만으로는 나중에 ack 할 수단이 없다. flush 는 **다음 리스너 진입 시점**에 그 파티션의 완료 연속 구간까지 수행한다 | `SPRK-ACKMODE-C2` + 위 "ack 호출 스레드" 행의 스레드 제약 |
|
||
| 유휴 파티션 flush | **`UNSUPPORTED_IMPL_DECISION`** — `concurrency` 를 파티션 수로 pin 하면(§3) 파티션이 pause 되거나 유휴인 동안 리스너 진입이 아예 없어 "다음 리스너 진입 시 flush" 가 발화하지 않는다. 컨테이너 idle 이벤트 계열 훅이 후보이나 **그 훅이 어느 스레드에서 발화하는지가 수집한 raw 에 없다**(잔여 `SOURCE_GAP-5`). trade-off: **리스너/consumer 스레드에서 실행되는 idle 훅만 사용**하고, 그런 훅이 없다고 확인되면 `MANUAL_IMMEDIATE` 대신 완료 즉시 ack 하는 인라인 모델(D4 대안)로 후퇴한다 — 별도 스케줄러 스레드에서 ack 를 호출하는 방식은 근거 없이 채택하지 않는다 | `SPRK-ACKMODE-C5`/`C6` (인접 API 의 스레드 구속), `SPRK-ACKMODE-C7` (concurrency 하향 조정) |
|
||
|
||
> **잔여 `UNSUPPORTED_IMPL_DECISION`**: `acknowledge()`(nack/부분배치가 아닌 일반 ack)를 **다른 스레드에서 호출했을 때의 동작**은 수집한 문서가 직접 규정하지 않는다(`SPRK-ACKMODE-C5`/`C6` 는 각각 `nack()`·부분 배치 ack 에 한정된 제약이다). trade-off: 두 인접 API 가 모두 스레드에 묶여 있으므로 **보수적으로 리스너 스레드 ack 를 계약으로 고정**한다 — 반대로 갔다가 틀리면 유실이지만, 이 방향으로 틀리면 성능만 손해다.
|
||
|
||
정상 경로 순서 (이 순서를 어기면 D3 위반):
|
||
|
||
워커 풀은 **이 leaf 전용**이며 파티션당 1 워커로 둔다 — sibling background-job 의 executor 를 공유하지 않는다(그 branch 의 saturation 정책이 레코드 거부를 뜻해 D5 와 충돌하므로). 컨테이너 스레드가 파티션 수만큼 생기므로(§3 `concurrency`) 총 스레드는 대략 `파티션 수 × 2` 이고, 워커 drain await 예산은 그 branch 의 shutdown 예산 안에 들어가야 한다(§엣지의 stop 계약).
|
||
|
||
```text
|
||
poll() → (파티션별 bounded queue 에 적재; 포화면 그 파티션만 pause — D5)
|
||
→ worker: application use case 실행
|
||
→ TransactionPort.inWrite { 비즈니스 write + InboxStorePort.insert } ← 여기서 커밋 (D10)
|
||
→ 커밋 성공 확인 후에만 offset ack ← (D3)
|
||
→ 큐 여유 생기면 resume — D5
|
||
```
|
||
|
||
| 규칙 | 근거 | 어겼을 때 |
|
||
|---|---|---|
|
||
| `enable.auto.commit=false` 를 **명시** 설정 (기본값이 `true` 이므로 안 끄면 자동 커밋됨) | `KAFKA-CONSCFG-C3` (기본 true), `KAFKA-CONSCFG-C4` (5000ms 주기) | 처리 완료와 무관하게 5초마다 커밋 → 유실 |
|
||
| poll 이 반환한 레코드를 큐에 넘기는 즉시 처리 완료로 간주하지 않음 | `KAFKA-OFFSET-C3` | 자동 커밋의 at-least-once 전제(`KAFKA-OFFSET-C2`)가 깨짐 |
|
||
| ack 는 DB 커밋 **이후**. 반대 순서(ack 먼저)는 금지 | `KAFKA-OFFSET-C4` | ack 후 커밋 실패 시 재전달 없이 영구 손실 = at-most-once 로 후퇴 |
|
||
| pause 는 poll 을 멈추는 것이 아니다 — 컨테이너는 계속 poll 하며 레코드만 안 가져온다 | `SPRK-PAUSE-C2` | poll 자체를 멈추면 `max.poll.interval.ms` 초과 → 그룹 이탈(§3) |
|
||
| **pause 단위는 파티션** — `pausePartition(TopicPartition)`/`resumePartition(TopicPartition)` 을 쓴다. 컨테이너 전역 `pause()` 는 한 파티션의 포화로 나머지 파티션까지 굶기므로 기본 경로가 아니다 | `SPRK-PAUSEPART-C1` (2.7~ 파티션 단위 API), `SPRK-PAUSEPART-C2` (poll 경계 반영) | 전역 pause 를 쓰면 D4 의 파티션별 독립 큐가 무의미해진다 |
|
||
| pause 요청과 실제 정지를 구분 — 파티션 단위는 `isPartitionPauseRequested()` ≠ `isPartitionPaused()` (컨테이너 단위의 `isPauseRequested()`/`isConsumerPaused()` 와 같은 구조) | `SPRK-PAUSEPART-C3`, `SPRK-PAUSE-C4` | 정지 전에 큐를 비었다고 판단해 resume → 포화 반복 |
|
||
| pause/resume 전이는 관측 가능해야 한다 — `ConsumerPartitionPausedEvent`/`ConsumerPartitionResumedEvent` 를 지표(D14 의 `consumer.paused.seconds`)로 연결 | `SPRK-PAUSEPART-C4` | 조용한 정체를 탐지할 방법이 없어진다 |
|
||
| poll 루프 스레드에서 블로킹 대기 금지 (큐 offer 는 non-blocking) | `SPRK-PAUSE-C2` 의 전제 + `KAFKA-CONSCFG-C1` | pause 여부와 무관하게 `max.poll.interval.ms` 타이머가 흐름 |
|
||
| **poll 배치 흡수 불변식** — pause 를 요청하는 high watermark 는 `큐 용량 − max.poll.records` 이하로 둔다. 즉 **직전 poll 이 반환한 배치를 통째로 넣을 여유가 남아 있을 때 pause 를 요청**한다 | `SPRK-PAUSEPART-C2` (파티션 pause 도 poll 경계에서 반영 — pause 요청 후 추가 유입이 **배치 1개로 상한**된다는 핵심 근거), `KAFKA-CONSCFG-C2` (`max.poll.records` 기본 500 = 흡수해야 할 최대치), `SPRK-PAUSE-C3` (컨테이너 레벨에서 확인된 보수적 상한 — 기본 pause 는 "이전 poll 의 모든 레코드 처리가 끝난 뒤" 발효. 이 옵션이 파티션 단위 API 에도 동일 적용되는지는 원문에 명시가 없어 **더 보수적인 쪽**으로 채택) | 불변식이 깨지면 "drop 금지(D5) · 블로킹 금지 · pause 미발효" 가 동시에 성립해 **합법적 행동이 남지 않는다** |
|
||
|
||
### 3. rebalance·poll 설정 계약 (D6)
|
||
|
||
> **Trace**: D6 (`KIP429-C1`~`C5`, `KAFKA-CONSCFG-C1`/`C2`/`C5`/`C6`/`C7`, `VGS-REBAL-C2`~`C5`).
|
||
>
|
||
> - **UNSUPPORTED_IMPL_DECISION**: (1) `max.poll.interval.ms`·`max.poll.records` 의 **실제 pin 값**. 공식 문서는 기본값만 말하고 VGS 사례의 값(600000ms / 5)은 100-consumer aiokafka 환경 전제라 그대로 옮길 수 없다. trade-off: 값을 지어내는 대신 **모델별 관계식**을 계약으로 두고 기본값은 공식 기본값을 상속한다 — **큐 기본 모델(D4 기본)**: `poll→enqueue 소요 + pending ack flush(동기 커밋) 소요 < max.poll.interval.ms`. poll 스레드가 핸들러를 실행하지 않으므로 핸들러 p99 는 이 식에 들어가지 않는다. 흡수 불변식(§2)은 이 식의 **대기 항을 0 으로 만드는 전제**이지 항이 아니다 — 여유가 확보돼 있으므로 enqueue 는 블로킹하지 않는다. **인라인 대안 모델(D4 대안)**: `핸들러 p99 × max.poll.records < max.poll.interval.ms`. (2) 큐 용량·high/low watermark 의 절대값 — 불변식(`high watermark ≤ 용량 − max.poll.records`)만 계약이고 수치는 설정으로 노출.
|
||
|
||
| 설정 | 공식 기본값 (근거) | 본 branch 의 계약 |
|
||
|---|---|---|
|
||
| `enable.auto.commit` | `true` (`KAFKA-CONSCFG-C3`) | **`false` 로 명시 pin** (D3) |
|
||
| `max.poll.interval.ms` | `300000` (`KAFKA-CONSCFG-C1`) | 실측 기반 명시 pin. 초과 시 consumer 가 실패로 간주되어 파티션이 재할당됨 |
|
||
| `max.poll.records` | `500` (`KAFKA-CONSCFG-C2`) | 큐 용량(§2)과 함께 결정. 배치 크기가 poll 주기 예산을 좌우 |
|
||
| `session.timeout.ms` / `heartbeat.interval.ms` | `45000` / `3000`, heartbeat 는 session 의 1/3 이하 권장 (`KAFKA-CONSCFG-C5`/`C6`) | 기본값 상속 — 본 branch 는 재정의하지 않음(처리 지연은 `max.poll.interval.ms` 축이 담당) |
|
||
| 컨테이너 `concurrency` | (프레임워크 속성) 파티션 수보다 크면 **하향 조정**된다 (`SPRK-ACKMODE-C7`) | 파티션 수를 상한으로 pin. D4 의 "병렬도 상한 = 파티션 수" 가 프레임워크 차원에서도 강제된다 — 큐 모델에서도 이 값을 넘겨 잡지 않는다 |
|
||
| `partition.assignment.strategy` | `[RangeAssignor, CooperativeStickyAssignor]` (`KAFKA-CONSCFG-C7`) | `CooperativeStickyAssignor` 로 pin. 기본 목록 덕에 `RangeAssignor` 만 제거하는 단일 rolling bounce 로 전환 가능 |
|
||
|
||
rebalance 리스너 계약 (D6):
|
||
|
||
| 콜백 | 계약 | 근거 |
|
||
|---|---|---|
|
||
| `onPartitionsRevoked` | 호출되면 그 파티션의 **완료분까지만** 커밋. **호출을 전제하지 않는다** — cooperative 에서는 아예 호출되지 않을 수 있다 | `KIP429-C5` |
|
||
| `onPartitionsAssigned` | 새로 배정된 파티션의 pause 상태를 초기화(resume)한다. 안 하면 배정받고도 소비하지 않는 좀비 파티션이 된다 — **`UNSUPPORTED_IMPL_DECISION`**: 재배정 시 파티션 pause 상태가 유지되는지 초기화되는지를 규정한 문서가 없다(`SPRK-PAUSEPART-C5` — 파티션 pause 페이지·컨테이너 pause 페이지 모두 rebalance 어휘 자체가 부재). trade-off: **보수적으로 명시 resume 을 호출**한다 — 이미 resume 상태에 resume 을 부르는 것은 무해하지만, pause 가 잔존하면 그 파티션이 조용히 멈춘다(비대칭 위험) | `SPRK-PAUSEPART-C5` (부재 확인), `SPRK-PAUSEPART-C3` (요청 vs 실제 정지 구분 API 로 상태 확인 가능) |
|
||
| `onPartitionsLost` | 이미 소유권을 잃은 뒤의 정리 전용 — 이 시점의 커밋은 무효로 간주 | `KIP429-C4` |
|
||
| 미완료 큐 항목 | 폐기(커밋하지 않음). 재할당 consumer 가 마지막 커밋 offset 부터 재소비 → **중복이지 유실 아님**, D10/D11 이 흡수 | `KAFKA-OFFSET-C4` + D10 |
|
||
|
||
### 4. 실패 분류 → 회수 경로 (D7, D8, D9)
|
||
|
||
> **Trace**: D7 (`SPRK-EHD-C1`~`C3`, `SPRK-ERRH-C1`, `CONF-POISON-C3`/`C4`) + D8 (`SPRK-RETRYTOPIC-C3`, `SPRK-ERRH-C2`, background-job D4 위임) + D9 (`SPRK-ERRH-C3`/`C4`/`C5`, modules.yaml 경계).
|
||
>
|
||
> - **UNSUPPORTED_IMPL_DECISION**: (1) **감사된 replay 의 기록 스키마**(누가/언제/어느 offset/결과) — Spring·Confluent·Uber 어느 문서도 replay audit 를 규정하지 않는다(3개 조사 모두 "외부 근거 부재" 로 확인). trade-off: 최소 필드(`replayedBy`, `replayedAt`, 원본 `topic/partition/offset`, `idempotencyKey`, 결과)를 inbox/dead-letter row 에 남기는 방향만 정하고 상세는 구현 시. (2) blocking retry 를 pause 형 backoff 로 전환하는 **실패율 임계** — 근거 없음. trade-off: backoff 총합이 `max.poll.interval.ms` 를 넘길 수 있으면 전환한다는 조건만 계약화. (3) **fatal 6종(`SPRK-ERRH-C1`) 외에 어떤 프로젝트 예외를 non-retryable 로 확장할지** — 프레임워크는 목록 확장 수단만 제공하고 무엇을 넣을지는 말하지 않는다. trade-off: 도메인 검증 실패 계열(`VALIDATION`/`DATA_INTEGRITY` category)을 우선 후보로 두되, 애매한 예외는 확장하지 않고 재시도 예산 소진에 맡기는 편이 유실보다 안전(D7 Open Risk 의 truncation 엣지도 이 원칙으로 흡수).
|
||
|
||
| 시나리오 | 분류 | 처리 경로 | offset | Trace |
|
||
|---|---|---|---|---|
|
||
| 역직렬화 실패 (깨진 바이트·스키마 불일치) | **non-retryable** | 리스너 호출 없이 error handler → dead-letter 회수 (D9) | 회수 성공 후 전진 | D7 (`SPRK-EHD-C2`/`C3`, `SPRK-ERRH-C1`) |
|
||
| handler/schema/version allowlist 위반 (파싱은 성공, 미등록 `(topic, eventType)`) | **non-retryable** | application 예외 → dead-letter 회수 | 회수 후 전진 | **D15**, D7 선택 조건 (§구현 가이드 6) |
|
||
| 다운스트림 일시 실패 (DB·외부 의존 타임아웃) | **retryable** | blocking bounded retry (backoff+jitter, max attempts 3 — background-job D4 위임) | 성공 시 전진 / 소진 시 아래 | D8 |
|
||
| 재시도 소진 | 종단 실패 | dead-letter 회수 + 감사 기록 | 전진 | D8, D9 |
|
||
| backoff 총합이 `max.poll.interval.ms` 를 넘길 위험 | — | 스레드 정지형 backoff 대신 **컨테이너 pause 형** backoff | — | D8 (`SPRK-ERRH-C2`) |
|
||
| dead-letter 회수 자체가 실패 | 종단 실패 | ack 하지 않음 → 재전달되어 재시도 (중복은 D10/D11 흡수) | 전진하지 않음 | D3, D9 |
|
||
|
||
dead-letter 발행 경계 (D9) — **모듈 규칙이 방식을 강제한다**:
|
||
|
||
- `DeadLetterPublishingRecoverer` 는 레코드를 보내기 위해 producer(`KafkaTemplate`/`KafkaOperations`)를 요구한다(`SPRK-ERRH-C4`). 그런데 inbound leaf 의 `allowed_dependencies` 에는 outbound leaf 가 없다(modules.yaml 실측).
|
||
- 따라서 **inbound leaf 가 producer 를 직접 들 수 없다**. dead-letter 발행은 `application-core` 에 정의한 outbound port 를 통해 나가고, 그 구현은 `adapter:outbound:messaging` 이 맡는다(Dependency Inversion). 대안으로 dead-letter 를 DB row 로만 표현하면 producer 자체가 필요 없다(D9 선택 조건).
|
||
- recoverer 를 명시 구성하지 않으면 기본 동작은 **로그만** 이고 dead-letter 로 가지 않는다(`SPRK-ERRH-C5`) — "설정 안 하면 안전" 이 아니라 "설정 안 하면 조용히 유실" 이다.
|
||
- 토픽 방식 채택 시 기본 명명은 `<originalTopic>-dlt`, 원본과 같은 partition 이며 DLT 토픽의 partition 수가 원본 이상이어야 한다(`SPRK-ERRH-C3`).
|
||
|
||
### 5. Inbox 스키마·dedupe key (D10, D11, D12, D13)
|
||
|
||
> **Trace**: D10 (`MSIO-IDEMPC-C2`~`C5`, ca-tmpl `TransactionPort`) + D11 (outbox D12/D14, `CLOUDEVT-C2`, `MSIO-IDEMPC-C4`) + D12 (#070 위임) + D13 (UNSUPPORTED — TTL 수치).
|
||
>
|
||
> - **UNSUPPORTED_IMPL_DECISION**: (1) 테이블·컬럼 물리 설계(DB 타입, 인덱스) — 근거 없음. trade-off: outbox `V3__outbox_event.sql` 선례의 컬럼 명명·인덱스 패턴을 따르는 편이 새 규칙을 만드는 것보다 일관적. (2) reaper 주기·보존 상한 — D13 의 관계식만 있고 수치 근거가 없다. trade-off: outbox reaper(`ca-skeleton.outbox.published-retention`) 를 모델로 설정 키로 노출.
|
||
|
||
| 항목 | 명세 | 근거 |
|
||
|---|---|---|
|
||
| 포트 | `InboxStorePort` (`application-core`, `dev/caskeleton/application/inbox/`) — **신규, 현재 부재** | D10, `planned` |
|
||
| 저장 어댑터 | `adapter:outbound:persistence-jpa` (`dev/caskeleton/adapter/outbound/persistence/`) — outbox 의 `OutboxStoreAdapter`/`OutboxEventEntity` 선례와 같은 자리. inbound leaf 는 포트만 호출하고 저장 구현을 알지 못한다 | D1(의존 제한), D10 |
|
||
| 마이그레이션 | 신규 inbox 테이블의 Flyway 스크립트. 버전 번호 배정·적용 순서·`out-of-order` 정책은 [[raw/branch-notes/feature-migration-startup-contract]] 소유 — 본 branch 는 소비자 | D10 (reference-only) |
|
||
| 트랜잭션 경계 | `TransactionPort.inWrite { 비즈니스 write + inbox insert }` — 포트 구현이 자체 트랜잭션을 열지 않는다(outbox `OutboxStorePort` 규약과 동일) | D10, ca-tmpl `TransactionPort` (`actually-implemented`) |
|
||
| 중복 탐지 | 유니크 제약 위반으로 INSERT 실패 → 트랜잭션 rollback → 중복 처리 원천 차단 | `MSIO-IDEMPC-C4` |
|
||
| dedupe key | `idempotencyKey`(= `eventId`, ULID) 1차 + 유니크 제약은 `(idempotencyKey, eventType)` 복합 | D11, outbox D12 |
|
||
| 금지 | `(topic, partition, offset)` 단독 key — outbox relay 재발행이 같은 논리 이벤트를 다른 offset 으로 싣는다 | D11 (outbox §엣지) |
|
||
| 기록 필드(최소) | `idempotencyKey`, `eventType`, `aggregateId`, `correlationId`, 처리 시각, 원본 `topic/partition/offset`(감사용) | D11, outbox D12 envelope |
|
||
| 대안 (별도 테이블 없이) | 비즈니스 엔티티 자체에 메시지 ID 저장 — 1 이벤트 = 1 row 인 경우만. 처리 이력 조회는 포기 | `MSIO-IDEMPC-C5`, D10 선택 조건 |
|
||
| owner token 재사용 | 기본 미사용(insert-once). worker fan-out 도입 시 #070 의 claim/lease 로 승격 | D12 |
|
||
| 보존(TTL) | 수치 미정. `retention.ms` 와 replay 창 중 긴 쪽 이상 + 무한 보관 금지 | D13 (UNSUPPORTED_DECISION) |
|
||
|
||
### 6. handler/schema/version allowlist (D15, D7)
|
||
|
||
> **Trace**: D15 (`SPRK-ERRH-C1` — 라우팅/시그니처 불일치 계열이 fatal, outbox D12 — envelope `eventType` required) + D7 (미등록 조합의 처리 경로를 공유).
|
||
>
|
||
> - **UNSUPPORTED_IMPL_DECISION**: (1) 등록부를 **코드에 둘지 설정에 둘지** — 외부 근거 없음. trade-off: 설정과 handler 가 갈라지면 "등록됐는데 handler 없음" 이 런타임에만 드러나므로 코드 등록부가 안전. (2) `schemaVersion` 축 — **현재 envelope 에 그 필드가 없다**(outbox D12 required 6필드 실측). trade-off: 지금 발명하지 않고 실효 키를 `(topic, eventType)` 2축으로 두되, 버전 축이 필요해지면 outbox D12 확장을 요청한다.
|
||
|
||
| 항목 | 명세 | Trace |
|
||
|---|---|---|
|
||
| allowlist 의 키 | `(topic, eventType)` — `schemaVersion` 은 envelope 확장 후 추가 | D15, outbox D12 |
|
||
| 2축이 실효를 갖는 전제 | 현재 내부 producer 는 `topic = eventType` 이라(#063 D5) 두 축이 **1:1 로 축약**되어 allowlist 가 걸러낼 것이 없다. 이 검사가 실제로 작동하는 경우는 (a) 외부 시스템이 우리 토픽에 발행하거나 (b) 한 토픽에 여러 `eventType` 을 싣는 매핑을 도입할 때다. 그 전까지 D15 의 실효 방어선은 **구독 목록 자체**(§위 행)이며, allowlist 는 그 시점을 대비한 계약이다 | D15, #063 D5 (`actually-implemented`) |
|
||
| 등록 주체 | 각 handler 가 자신이 담당하는 조합을 선언. 별도 yaml/env 등록부를 두지 않는다 | D15 |
|
||
| 구독 토픽 목록의 출처 | **handler 선언의 topic 합집합** — 별도 env/yaml 로 토픽을 나열하지 않는다. 따라서 "구독했는데 handler 없음" 은 구조적으로 발생하지 않고, 반대로 handler 가 늘면 구독도 함께 는다 | D15 |
|
||
| 기동 시 검증 | 같은 `(topic, eventType)` 조합을 두 handler 가 선언(중복)하면 기동 거부. handler 가 0개면(= 구독 토픽 0개) 이 leaf 자체가 비활성으로 취급되어 listener container 를 만들지 않는다(D2 의 활성화 게이트와 동일 원칙). 기동 거부의 실패 표현은 [[raw/branch-notes/feature-migration-startup-contract]] 계약을 따른다 | D15, D2 (reference-only) |
|
||
| 미등록 조합 수신 | 처리하지 않고 **non-retryable** 로 분류해 D9 회수 경로. 조용한 skip 금지 | D15, D7 |
|
||
| 검사 위치 | 역직렬화 성공 **이후**, use case 호출 **이전** — D7 의 deserializer 경계 검사와 단계가 다르다(그쪽은 파싱 자체의 실패) | D15, D7 |
|
||
|
||
### 7. 계약 값 — 전부 신규 제안 (D14)
|
||
|
||
> **Trace**: D14 (ca-tmpl `docs/registries/*.yaml` 실측 — consumer/inbox row 부재, `Category.java` 10-value enum).
|
||
>
|
||
> - **UNSUPPORTED_IMPL_DECISION**: 아래 code/metric/env 의 **명명 자체**는 외부 근거가 없다. trace-off: 기존 registry 의 도메인 접두 관행(`OUTBOX_*`/`JOB_*`, `outbox.*`/`job.*`, `APP_MESSAGING_*`)을 그대로 따르는 편이 새 네이밍 축을 만드는 것보다 일관적. **아래는 제안이며 registry 반영 전까지 기존 값처럼 인용 금지.**
|
||
|
||
| 종류 | 제안 값 | 제안 속성 | 대응 결정 |
|
||
|---|---|---|---|
|
||
| error code (신규 제안) | `CONSUMER_DESERIALIZATION_FAILED` | category `DATA_INTEGRITY` 또는 `PERMANENT_DEPENDENCY` 중 택일(둘 다 기존 enum 값), `retryable: false`, runbook 필수 | D7, D14 |
|
||
| error code (신규 제안) | `CONSUMER_DEAD_LETTER` | category `INTERNAL`, `retryable: false` — `OUTBOX_DEAD_LETTER`/`JOB_DEAD_LETTER` 선례와 동형 | D9, D14 |
|
||
| error code (신규 제안) | `INBOX_DUPLICATE_SKIPPED` | 오류가 아니라 정상 경로 — **code 대신 metric 으로만 표현**하는 편이 registry 오염이 적다(대안 명시) | D10, D14 |
|
||
| metric (신규 제안) | `consumer.records.total{topic,outcome}` | outcome ∈ {PROCESSED, DUPLICATE, RETRIED, DEAD} — `job.retry.total` 의 tag 패턴 참고 | D7~D11 |
|
||
| metric (신규 제안) | `consumer.lag` / `consumer.queue.depth` / `consumer.paused.seconds` | pause 지속·큐 적체가 조용한 정체의 유일한 관측 지점 | D5 |
|
||
| env key (신규 제안) | `APP_MESSAGING_KAFKA_CONSUMER_*` (group id, 동시성, 큐 용량, `max.poll.*`) | 기존 `APP_MESSAGING_KAFKA_BROKERS` 접두 관행 상속. env key owner 는 `feature-env-driven-runtime-configuration`/`feature-integration-adapter-templates` — 등록은 협의 필요 | D2, D6 |
|
||
| env key (**본 branch 대상 아님**) | `security.protocol`·TLS/SASL 자격증명 계열 | broker 접속 보안 키는 [[raw/branch-notes/feature-kafka-producer-runtime-contract]] **D7** 소유 — consumer 는 같은 키 집합을 상속하고 신규 정의하지 않는다 | D2 |
|
||
| runbook (의무) | `runbook://consumer/dead-letter`, `runbook://consumer/poison-record` | registry 규약상 `retryable=true` 행과 `retryable=false` + INTERNAL 계열은 runbook_link 필수 | D9, D14 |
|
||
|
||
<!-- section-id: edge-failure-dependency -->
|
||
## 엣지·실패·의존
|
||
|
||
- **실패·엣지 경로**:
|
||
- **DB 커밋 성공 후 ack 전 crash** → 같은 메시지 재전달. 기대 동작: inbox 유니크 제약 위반으로 중복 감지 → 처리 없이 ack (D10, D11). at-least-once 의 정상적 결과다.
|
||
- **ack 후 DB 커밋 실패** → 발생해서는 안 되는 순서. D3 이 금지하는 배치이며, 어기면 재전달 없는 영구 손실(`KAFKA-OFFSET-C4`).
|
||
- **poll 루프 스레드 블로킹(GC·큐 offer 대기)** → pause 여부와 무관하게 `max.poll.interval.ms` 초과 → 그룹 이탈·파티션 재할당(`KAFKA-CONSCFG-C1`). 기대 동작: 큐 offer 는 non-blocking, 포화는 pause 로 표현(D5, §구현 가이드 2).
|
||
- **rebalance 중 in-flight 레코드** → 미완료분은 커밋하지 않고 폐기. 재할당 consumer 가 마지막 커밋 offset 부터 재소비하여 중복 발생 — D10/D11 이 흡수. `onPartitionsRevoked` 가 호출되지 않을 수 있으므로 그 콜백을 유일 체크포인트로 삼지 않는다(`KIP429-C5`).
|
||
- **재배정 후 pause 상태 미초기화** → 파티션을 배정받고도 fetch 하지 않는 좀비 파티션. 기대 동작: `onPartitionsAssigned` 에서 강제 resume(§구현 가이드 3).
|
||
- **poison record 무한 재시도** → offset 이 전진하지 않아 해당 파티션 소비가 정지(`CONF-POISON-C3`/`C4`). 기대 동작: deserializer 경계에서 non-retryable 로 분류해 첫 실패에 회수(D7).
|
||
- **dead-letter recoverer 미구성** → 재시도 소진 후 기본 동작이 로그만이라 조용히 유실(`SPRK-ERRH-C5`). 기대 동작: recoverer(또는 DB dead-letter row) 구성 없이 이 leaf 를 활성화하지 못하게 기동 검증.
|
||
- **producer 신뢰 붕괴 (서로 다른 논리 이벤트가 같은 `idempotencyKey`)** → 정상 이벤트를 중복으로 오판해 **조용히 누락**. 기대 동작: `(idempotencyKey, eventType)` 복합 제약이 불일치를 탐지해 경고로 격상(D11).
|
||
- **inbox TTL < replay 창** → 감사된 replay 로 되돌린 메시지가 "처음 보는 메시지" 로 재처리됨(D13 관계식이 막으려는 모순).
|
||
- **backoff 총합 > `max.poll.interval.ms`** → 재시도 도중 그룹 이탈. 기대 동작: 스레드 정지형이 아니라 컨테이너 pause 형 backoff 로 전환(D8, `SPRK-ERRH-C2`).
|
||
- **배포 순서 편차 (producer 선배포 → 구 consumer 가 신규 이벤트 타입을 모름)** — 현재 producer 는 `topic = eventType` 으로 발행하므로(#063 D5, `actually-implemented`) 이 경우가 **두 갈래로 갈린다**:
|
||
- (a) **미구독 신규 토픽**: 구독 목록이 handler 선언의 합집합이므로(§구현 가이드 6) 그 토픽을 아무도 구독하지 않는다 → 레코드가 consumer 에 **도달조차 하지 않고** 토픽에 적체되다 `retention.ms` 만료로 유실될 수 있다. D15 의 allowlist 는 이 경로를 막지 못한다. 기대 동작: consumer 배포 전까지 적체를 견디도록 해당 토픽 retention 을 확보하고, 미구독 토픽 존재를 운영이 인지할 수단(브로커 측 토픽 목록 대조)이 필요하다 — **본 branch 결정 범위 밖의 운영 절차**.
|
||
- (b) **구독 중인 토픽에 미등록 `eventType` 투입**(외부 producer 또는 1토픽-다eventType 매핑을 쓰는 경우): D15 에 따라 non-retryable 로 분류되어 dead-letter 로 회수된다 → 유실이 아니며 consumer 배포 후 감사된 replay(D9)로 재처리. 관측은 `consumer.records.total{outcome=DEAD}` 급증(D14 제안 지표).
|
||
- **정지(SIGTERM) 시점의 in-flight** → 파티션 큐에 남은 레코드 · 아직 flush 되지 않은 pending ack · `TransactionPort.inWrite` 실행 중인 워커가 동시에 존재한다. 기대 동작(**stop 계약**): ① 모든 파티션 `pausePartition` → ② 워커 drain 을 bounded 하게 await → ③ **완료 연속 구간까지 pending ack flush** → ④ consumer close. 예산 안에 끝나지 않은 미완료분은 **ack 하지 않고 폐기**한다 — 재시작 후 재전달되며 중복은 D10/D11 이 흡수한다(유실보다 중복을 택하는 D3 와 같은 방향). 정지 순서상의 위치와 예산 배분은 [[raw/branch-notes/feature-runtime-health-lifecycle-contract]] **D4**(SIGTERM → readiness DOWN → inflight drain → outbound 컴포넌트 descending stop → exit)와 [[raw/branch-notes/feature-background-job-async-contract]](executor `setWaitForTasksToCompleteOnShutdown(true)` + `setAwaitTerminationSeconds(19)`, 기본값은 즉시 interrupt)의 예산 안에 들어가야 하므로 **두 owner 와 phase 배치 협의 필요**. F2 의 지연 ack 때문에 이 경로는 매 배포마다 반드시 발생한다.
|
||
- **유휴 파티션의 pending ack 미발화** → 워커가 완료 offset 을 적재했는데 리스너 진입이 끊겨 ack 이 나가지 않음. 기대 동작: 유휴 flush 경로가 이를 밀어낸다(§구현 가이드 2 "ack 발화 지점"). 없으면 재기동 시 대량 재전달 + `consumer.lag` 왜곡.
|
||
- **broker 장기 다운** → 큐 적체 → pause 지속. 관측 지점이 없으면 조용한 정체가 된다 — `consumer.paused.seconds`/`consumer.queue.depth`(D14 신규 제안)가 필요한 이유.
|
||
- **다른 계약 의존**:
|
||
- [[raw/branch-notes/feature-kafka-producer-runtime-contract]] (`WI-CA-SKELETON-OPERATIONAL-CONTRACT-063`, **D1~D13 작성 완료 — 2026-07-28 확인**):
|
||
- **D5**(순서 보장은 파티션 단위, `key = aggregateId` 로 per-aggregate FIFO 대응) — 본 branch D4 의 전제를 producer 측에서 실현한다. 그 노트가 `OutboxMessagePublishAdapter` 의 `topic=eventType, key=aggregateId` 를 `actually-implemented` 로 확인했으므로 D4 의 "파티션 = aggregate 단위" 가정은 근거를 얻는다. 다만 **그 D5 의 Open Risk(파티션 수를 늘리면 같은 `aggregateId` 가 다른 파티션으로 가서 per-aggregate FIFO 가 깨짐)가 본 branch 의 순서 계약에도 그대로 전이**된다.
|
||
- **D2**(seam 유지 + 스켈레톤이 spring-kafka 기반 기본 구현 제공 + broker 미선택 시 auto-config 비활성) — 본 branch D2 가 같은 축으로 정렬. 같은 D2 가 "모듈 registry migration(19→20)은 #064 소유" 로 본 branch D1 을 명시 승인한다.
|
||
- **D10**(`OutboundMessage` 에 headers 추가 — `mdc-keys` 의 `propagation: [message]` 4종 전파) — consume 시 복원해야 할 header 집합의 producer 측 계약. 그 D10 이 바뀌면 본 branch 의 MDC 복원 대상이 바뀐다.
|
||
- **D3**(`MessageBroker.send` 반환 확장) / **D11**(producer 전용 error code 미생성) — 본 branch D14 가 "신규 코드 최소화" 방향을 참고할 선례.
|
||
- [[raw/branch-notes/feature-idempotency-ownership-protocol-contract]] (`WI-CA-SKELETON-OPERATIONAL-CONTRACT-070`) — D12 의 owner token 재사용 판단이 이 branch 의 claim/lease API 확정에 의존. **현재 스캐폴딩(D-row 0개, 2026-07-28 확인)** — 확정 후 D12 재검토.
|
||
- [[raw/branch-notes/feature-capability-provider-selection-contract]] **D13**(비활성 capability 는 연결·워커·스키마·health contributor 를 만들지 않는다) — D2 가 SDK 를 classpath 에 올리는 순간 이 계약과 충돌할 수 있다. consumer 는 auto-config 활성화가 곧 broker 연결 시도이므로 producer 보다 충돌이 즉각적이다.
|
||
- [[raw/branch-notes/feature-domain-event-outbox-contract]] D7·D12·D14 — envelope 필수 필드와 `idempotencyKey` scope 를 consume 한다(D11). envelope 가 바뀌면 dedupe key 계약이 연동 변경된다.
|
||
- [[raw/branch-notes/feature-background-job-async-contract]] D4 — backoff/max attempts/DLQ 어휘를 consume 한다(D8). 단 **saturation 정책(`AbortPolicy`)은 consume 하지 않는다** — 거부는 레코드 유실이라 at-least-once 를 깨므로 consumer 경계에서는 pause 로 대체(D5).
|
||
- [[raw/branch-notes/feature-integration-adapter-templates]] — `APP_MESSAGING_BROKER`/`APP_MESSAGING_KAFKA_BROKERS` env key owner. consumer 용 신규 키(D14) 등록 시 협의 대상.
|
||
- [[raw/branch-notes/feature-operational-error-observability-foundation]] — `correlation_id` 의미·전파 SSOT(`mdc-keys.yaml` 의 `propagation` 에 `message` 포함). consume 시 envelope `correlationId` → MDC 복원 의무는 이 계약을 따른다(reference-only).
|
||
- [[raw/branch-notes/feature-schema-serialization-contract]] — payload **직렬화·스키마 진화** 정책 owner (Avro/JSON, date/decimal 표현 등). ⚠️ **handler/schema/version allowlist 는 이 branch 가 소유하지 않는다** — 초안에서 그쪽으로 위임한다고 적었으나 2026-07-28 grep 결과 handler/topic/routing 언급 0건으로 확인돼 본 branch 의 D15 로 회수했다(§Audit `FALSE_DELEGATION`). 이 계약과의 실제 접점은 payload 직렬화 형식뿐이다.
|
||
- [[raw/branch-notes/feature-runtime-health-lifecycle-contract]] **D4** — SIGTERM 이후의 정지 순서(readiness DOWN → inflight drain → outbound 컴포넌트 descending stop → exit) owner. 본 branch 의 stop 계약(§엣지의 pause→drain→pending ack flush→close)이 그 순서의 어느 phase 에 들어가는지 **협의 대상**이다. 그 D4 가 바뀌면 본 branch 의 정지 시퀀스가 연동 변경된다.
|
||
- [[raw/branch-notes/feature-background-job-async-contract]] — retry 어휘(D8) 외에 **shutdown 예산**도 이 계약을 따른다: executor `setWaitForTasksToCompleteOnShutdown(true)` + `setAwaitTerminationSeconds(19)`(기본은 즉시 interrupt). 워커 drain await 예산이 이 안에 들어가야 한다. **워커 풀 소유는 본 branch** — 그 branch 의 saturation 정책(`AbortPolicy`)은 consume 하지 않는다(D5).
|
||
- [[raw/branch-notes/feature-migration-startup-contract]] — 신규 inbox 테이블의 Flyway 마이그레이션 순서·번호 배정과, "필수 구성 누락 시 기동 거부" 를 어떤 실패로 표현할지(§엣지의 recoverer 미구성 기동 검증)의 owner. 본 branch 는 소비자다.
|
||
- [[raw/branch-notes/feature-kafka-producer-runtime-contract]] **D7** (`security.protocol` 명시 선택 + prod 에서 `PLAINTEXT` 기동 거부 + TLS/SASL 자격증명은 `secrets-classification.yaml` 의 secret tier) — consumer 도 **같은 broker 접속 계약을 상속**한다. 보안 관련 env key 는 그 branch 소유이며 본 branch 가 신규 정의하지 않는다(D14 의 신규 제안 대상 밖).
|
||
|
||
<!-- section-id: claims-to-verify -->
|
||
## 검증해야 할 주장 / Claims To Verify
|
||
|
||
| Claim | Why uncertain | How to verify | Status |
|
||
|---|---|---|---|
|
||
| 본 branch 가 소유할 관심사가 sibling branch 결정과 겹치지 않는다 | 스캐폴딩 시점에는 D-row 가 없어 경계가 문장으로만 존재 | `/branch-spec` 후 `/sync` 실행 — owner 중복 검출 | `needs-confirmation` |
|
||
| "파티션 하나는 consumer group 안에서 정확히 한 consumer 가 소비한다" (D4 의 전제) | 현재 근거가 회사 블로그 1건(`UBER-REPROC-C2`)뿐 — Kafka 공식 verbatim 미수집. `kafka.apache.org/documentation` 은 JS SPA 라 정적 fetch 불가가 이미 확인됨 | Kafka 공식 정적 페이지(예: `ConsumerConfig`/`KafkaConsumer` javadoc 의 group management 절) 또는 Confluent 미러에서 verbatim 수집 후 D4 근거 격상 | `needs-confirmation` |
|
||
| 파티션 단위 pause/resume API(`pausePartition`/`resumePartition`)가 D4·D5 조합(파티션별 큐 + 그 파티션만 pause)을 지원한다 | 2026-07-28 `[[raw/official-docs/spring-kafka-pause-resume-partitions-on-listener-containers]]` 수집으로 API 존재·타이밍·상태조회는 확인됨(`SPRK-PAUSEPART-C1`~`C4`). **단, rebalance 재배정 시 이 pause 상태가 유지되는지 초기화되는지는 이 문서도 다루지 않음**(`SPRK-PAUSEPART-C5` — 부재 확인) | 해소(API 존재) — 잔여: Spring Kafka 소스 코드(`KafkaMessageListenerContainer`) 또는 통합 테스트로 재배정 시 pause 상태 동작 검증 | `needs-confirmation` (API 존재는 confirmed, rebalance 상호작용은 미확인 유지) |
|
||
| SPI seam(D2)만으로 D3·D5·D7·D8·D9 정책이 fork 구현에서 실제로 지켜지는지 검증 가능하다 | 정책의 실행 주체가 fork 의 client 구현이므로, 스켈레톤이 계약 테스트를 어떻게 제공할지 미확정 | fake client seam 기반 계약 테스트 스위트를 설계해, ack 순서·pause 전이·poison 회수·dead-letter 경로를 fake 로 단언 가능한지 실증 | `planned` |
|
||
| inbox 유니크 제약 위반이 중복 처리를 실제로 rollback 시킨다 (D10) | `MSIO-IDEMPC-C4` 는 패턴 카탈로그의 서술이며 ca-tmpl 의 `TransactionPort` + JPA 조합에서의 동작은 별도 검증 필요 | 계약 테스트: 동일 `idempotencyKey` 메시지 5회 전달 → 비즈니스 row 1건, inbox row 1건, 처리 횟수 1회 단언 | `planned` |
|
||
| ack 가 DB 커밋 이후에만 발생한다 (D3) | 코드 순서를 지키는지는 리뷰로 보장되지 않는다 | 계약 테스트: 커밋 직전 예외 주입 → 오프셋이 전진하지 않고 재전달됨을 단언 / 커밋 성공 후 ack 예외 주입 → 재전달 시 중복 스킵됨을 단언 | `planned` |
|
||
| rebalance 중 재할당 파티션이 pause 상태로 남지 않는다 (D6) | `onPartitionsRevoked` 미호출 가능성(`KIP429-C5`)과 결합하면 상태 초기화 누락이 조용히 남는다 | 계약 테스트: pause 상태에서 파티션 재할당 시뮬레이션 → `onPartitionsAssigned` 후 해당 파티션이 resume 상태임을 단언 | `planned` |
|
||
| dead-letter 발행 경로가 모듈 경계를 위반하지 않는다 (D9) | producer 요구(`SPRK-ERRH-C4`)와 inbound leaf 의존 제한이 충돌하므로 배선 실수가 나기 쉽다 | `./gradlew verifyCleanArchitectureDependencies` + `:app-bootstrap:test --tests '*CleanArchitectureTest'` 통과 확인 (inbound leaf 가 outbound leaf 를 import 하면 실패) | `planned` |
|
||
| producer 가 동일 논리 이벤트에 항상 같은 `eventId` 를 재사용한다 (D11 의 전제) | outbox branch 에 이 보장을 명시한 D-row 가 없음(2026-07-28 확인) | outbox relay 재발행 시 `eventId` 재사용 여부를 `PublishPendingOutboxEventsUseCase` 코드로 확인하고, 필요하면 outbox branch 에 D-row 추가 요청 | `needs-confirmation` |
|
||
|
||
## 관심사 커버리지 (coverage-auditor 자동 생성 — 있을 때)
|
||
|
||
> 2026-07-28 `coverage-auditor` 1회차 결과 + 그 지적을 반영한 loop 1 수정 상태. governing doc: [[raw/project-notes/ca-skeleton-operational-contract]] (§8.0 `WI-CA-SKELETON-OPERATIONAL-CONTRACT-064`, §Owner Map) + 분해 설계 §4.2 L146. **1회차 판정은 `Not-covered`(Blocking 1 — allowlist)** 였고, 아래 표는 D15 신설 후 상태다.
|
||
|
||
| 관심사 | 상태 | owner | 심각도 | 근거 |
|
||
|--------|------|-------|--------|------|
|
||
| 신규 inbound leaf 등록 + 모듈 registry migration(19→20) | covered-here | — | — | D1 (`.harness/project/modules.yaml` 실측, #063 D2 가 소유권 명시 승인) |
|
||
| manual acknowledgement (use case 성공 + DB 커밋 이후 ack) | covered-here | — | — | D3 |
|
||
| handler/schema/version allowlist | covered-here | — | — | **D15 (2026-07-28 loop 1 신설)** — 1회차에서 `missing` 🔴 Blocking 이었고, 초안의 위임 주장이 거짓으로 확인돼(§Audit `FALSE_DELEGATION`) 본 branch 로 회수 |
|
||
| bounded concurrency·queue + pause/resume backpressure | covered-here | — | — | D4, D5 |
|
||
| rebalance·`max.poll` 처리 + poison/역직렬화 실패 분류 | covered-here | — | — | D6, D7 |
|
||
| 재시도 + DLT + 감사된 replay | covered-here | — | — | D8, D9 (replay audit 스키마는 `UNSUPPORTED_IMPL_DECISION` — coverage gap 아님) |
|
||
| `InboxStorePort` 트랜잭션 커밋 규칙 + dedupe key | covered-here | — | — | D10, D11 |
|
||
| 상속: at-least-once + 멱등 consumer/inbox, exactly-once 미주장 | covered-here | — | — | 상속 표 branch application + D3·D10·D11 |
|
||
| 상속: owner token 재사용 여부 | delegated | [[raw/branch-notes/feature-idempotency-ownership-protocol-contract]] | 🟡 Should-fix | D12 — 위임 링크는 있으나 owner 가 스캐폴딩(D-row 0개). `needs-approval` + TODO 로 추적 중 |
|
||
| consume 시 correlationId → MDC 복원 | delegated | [[raw/branch-notes/feature-operational-error-observability-foundation]] | OK | §엣지·의존 링크. owner 는 `verified` 상태 |
|
||
| broker 접속 보안(`security.protocol`·TLS/SASL secret) | delegated | [[raw/branch-notes/feature-kafka-producer-runtime-contract]] D7 | OK | §엣지·의존 + §구현 가이드 7 (본 branch 신규 정의 없음) |
|
||
| inbox 테이블 마이그레이션·기동 검증 표현 | delegated | [[raw/branch-notes/feature-migration-startup-contract]] | OK | §엣지·의존 + §구현 가이드 5 |
|
||
| consumer/inbox 신규 error code·metric·env key 제안 | covered-here | — | — | D14 (registry 부재 확인 후 "신규 제안" 라벨) |
|
||
|
||
## Audit & Findings (2026-07-28 `/branch-spec` ground-truth 대조)
|
||
|
||
> ca-tmpl 코드·registry 실측과 인용 재검증에서 나온 사항. 자동 수정하지 않고 기록만 한다.
|
||
|
||
| Finding | 분류 | 내용 | 조치 |
|
||
|---|---|---|---|
|
||
| consumer·inbox 인프라 전부 부재 | `IMPLEMENTATION_GAP` | 2026-07-28 `src/` grep: `InboxStorePort` 0건, inbound messaging leaf 없음(`src/adapter/inbound/` = web/grpc/graphql/websocket 4종), modules.yaml 19개. 존재하는 것은 outbound 측 `MessageBroker`/`KafkaSender`/`KafkaMessageBroker`/`MessagingConfig` 와 `TransactionPort` 뿐 | §구현 가이드의 신규 항목을 전부 `planned` 로 표기. 구현 branch 착수 시 해소 |
|
||
| SPI 대칭성이 consumer 측에서 깨질 수 있음 | `SPI_ASYMMETRY` | outbound 코드는 현재 SDK-free seam(`KafkaSender`)으로 성립하지만, consumer 측 정책(D7·D8·D9)의 근거는 전부 `spring-kafka` API(`ErrorHandlingDeserializer`/`DefaultErrorHandler`/`DeadLetterPublishingRecoverer`)다. seam 만 두면 정책의 *모양*만 규정되고 실행은 fork 몫이 된다 | sibling #063 D2 가 "스켈레톤이 spring-kafka 기반 기본 구현 제공" 으로 이미 축을 정했으므로 본 branch D2 를 그 축에 정렬해 해소. 잔여 위험(classpath 오염 ↔ capability-provider D13)은 D2 Open Risk 로 이월 |
|
||
| 본 노트의 1차 D2 가 현행 코드만 보고 작성돼 sibling 결정과 어긋났음 | `SIBLING_DRIFT` | 세션 초반 확인 시 #063 은 스캐폴딩(192줄, D-row 0)이었으나 같은 날 D1~D13 이 작성됨(487줄). 그 D2 가 "skeleton 이 spring-kafka 기반 기본 구현을 제공" 을 결정해, 코드 실측(`KafkaSender` javadoc "The skeleton carries no Kafka SDK dependency")만으로 세운 본 노트의 초안 D2(순수 SPI)와 충돌 | 본 D2 를 #063 D2 축으로 **재작성 완료**(2026-07-28). 두 D2 모두 `proposed` 이므로, SDK 를 classpath 에 올릴지 여부는 두 branch 가 **함께** 확정해야 한다 |
|
||
| dead-letter 발행이 모듈 경계와 충돌 | `MODULE_BOUNDARY` | DLT 발행에 producer 필수(`SPRK-ERRH-C4`) vs inbound leaf 의 `allowed_dependencies` 에 outbound leaf 없음(modules.yaml 실측) | D9 에서 application-core port 경유로 해소. 대안(DB dead-letter row)도 명시 |
|
||
| 파티션-소비자 배타 배정의 공식 근거 미수집 | `SOURCE_GAP-1` | D4 의 핵심 전제가 회사 블로그 1건에만 의존. `kafka.apache.org/documentation` 은 JS SPA 로 정적 fetch 불가(선례: `raw/official-docs/kafka-message-delivery-semantics-design.md` §URL Fetch 실패 기록) | §검증해야 할 주장에 등재. 공식 정적 페이지에서 verbatim 수집 후 D4 격상 |
|
||
| spring-kafka 층의 ack 메커니즘 근거 미수집 | `SOURCE_GAP-5` | D2 가 프레임워크를 고정했는데 vault 에 `AckMode`/`Acknowledgment` 를 다루는 raw 가 0건이었다 | **대부분 해소**(2026-07-28 loop 1) — `raw/official-docs/spring-kafka-ack-mode-manual-commit-and-concurrency.md` 수집(`SPRK-ACKMODE-C1`~`C7`). ack 모드·기본값 `BATCH`·리스너 타입 제약·concurrency 하향 조정 확보. **잔여**: 일반 `acknowledge()` 의 호출 스레드 규칙과 ack 순서 제약 문장은 그 페이지에 **부재 확인**(인접 `nack()`·부분배치 제약만 존재) → §구현 가이드 2 의 잔여 `UNSUPPORTED_IMPL_DECISION` 으로 라벨링. 인접 페이지 "Manually Committing Offsets" 재조사 후보 |
|
||
| 파티션 단위 pause API 근거 미수집 | `SOURCE_GAP-2` | 1차 수집분(컨테이너 레벨 pause/resume 페이지)에 `pausePartition`/`resumePartition` 부재(`SPRK-PAUSE-C5`) | **해소**(2026-07-28 loop 1) — `raw/official-docs/spring-kafka-pause-resume-partitions-on-listener-containers.md` 수집(`SPRK-PAUSEPART-C1`~`C5`). D5 근거 보강 완료. 단 **rebalance 시 pause 상태의 운명은 그 문서에도 없음이 확인**돼(`C5`) §구현 가이드 3 의 해당 행은 `UNSUPPORTED_IMPL_DECISION` 로 라벨링 |
|
||
| 감사된 replay 의 외부 근거 부재 | `SOURCE_GAP-3` | Spring / Confluent / Uber / AWS 어느 문서도 "누가·언제·무엇을 재처리했는가" 의 audit trail 을 규정하지 않음(3개 조사에서 각각 확인) | §구현 가이드 4 의 `UNSUPPORTED_IMPL_DECISION` 로 라벨링. 최소 필드만 방향 제시 |
|
||
| inbox TTL 수치의 외부 근거 부재 | `SOURCE_GAP-4` | CloudEvents / AWS Prescriptive Guidance / microservices.io 어디에도 수치 없음. 실무 관행값(~72시간, 15일)은 API-level idempotency-key 도메인 사례라 전용 불가 | D13 을 `UNSUPPORTED_DECISION` 으로 라벨링하고 관계식만 고정 |
|
||
| consumer/inbox 계약 값 registry 미등록 | `REGISTRY_GAP` | `error-codes.yaml`·`metrics.yaml`·`env-keys.yaml` 에 consumer/inbox row 0건 (2026-07-28 grep) | D14 로 "신규 제안" 표기. 구현 branch 가 `owner_branch` 를 본 branch 로 등록 |
|
||
| 게이트 루프 천장에서 종료 — depth 미통과 상태 | `GATE_CEILING` | `/branch-spec` 의 루프 천장(2회)에 도달했다. **coverage 는 `Covered`(Blocking 0) 로 통과**했으나 **depth 는 마지막 감사 시점에 `Not ready`(Blocking 1 / Should-fix 4 / Advisory 2)** 였다. Blocking 은 "정지(SIGTERM) 시점 in-flight 처리 계약 부재" 였고, 감사기가 제시한 처방(owner 링크 2건 + 엣지 1행 + stop 계약 1행, 새 조사 불필요)을 **감사 이후에 적용**했다 — 즉 **이 수정은 재감사로 검증되지 않았다** | 다음 세션에서 `/depth feature-kafka-consumer-inbox-contract` 를 먼저 재실행해 Blocking 해소를 확인할 것. 함께 적용한 Should-fix 4건(ack 핸들 보관·유휴 flush 라벨 / poll 예산 식 교정 / 배포 편차 엣지 2갈래 분리 + D15 2축 전제 / 워커 풀 소유·자원 상한)도 같은 재감사에서 확인 대상. **미적용 잔여**: `SOURCE_GAP-1`(D4 의 "파티션당 1 consumer" official 근거 — 기존 KafkaConsumer Javadoc 에서 claim 추가 추출로 닫힘) |
|
||
| 초안이 존재하지 않는 위임처를 가리킴 | `FALSE_DELEGATION` | 초안 §엣지·의존이 "handler/schema/version allowlist 의 형식은 [[raw/branch-notes/feature-schema-serialization-contract]] 를 따른다" 고 적었으나, 2026-07-28 grep 결과 그 branch 는 handler/topic/routing 을 **한 번도 언급하지 않는다**(Avro/JSON·date/decimal 직렬화 전용). governing 설계 §4.2 L146 과 본 노트 §포함 범위가 모두 이 관심사를 **본 branch 것**으로 명시한다 | **D15 신설로 회수**(2026-07-28 loop 1). §엣지·의존의 위임 문장도 정정 — 그 계약과의 실제 접점은 payload 직렬화 형식뿐임을 명시 |
|
||
| 동시성 모델의 선택 축이 성립하지 않았음 | `INCOHERENT_CRITERION` | 초안 D4 가 "파티션 수 이내면 단일 스레드, 그 이상이면 큐로 확장" 이라 썼는데, D4 자신이 파티션당 직렬을 못박으므로 두 형태의 병렬도 상한이 동일하다 — 확장 트리거가 성립 불가. 반면 §구현 가이드 2 의 시퀀스와 D5 의 "항상 pause" 는 큐 모델을 무조건 전제 | **D4 재작성**(2026-07-28 loop 1) — 축을 "poll 스레드를 처리 지연에서 분리할 필요가 있는가" 로 교체하고 큐 모델을 스켈레톤 기본으로 고정, 인라인 처리를 조건부 대안으로 강등 |
|
||
| 의존 sibling 1종이 스캐폴딩 상태 | `DEPENDENCY_SCAFFOLD` | `depends_on` 중 **#063(producer)은 D1~D13 작성 완료**(2026-07-28 확인)이나 **#070(idempotency owner token)은 여전히 D-row 0개** — D12 의 재사용 판단 근거가 아직 문장 수준 | D12 를 `needs-approval` 로 유지. #070 의 `/branch-spec` 완료 후 재검토. #063 쪽 의존은 §엣지·실패·의존에 D-row 단위로 명시 완료 |
|
||
|
||
## 마주친 문제
|
||
|
||
아직 없음.
|
||
|
||
## 묶음 (이 branch에서 파생된 자료)
|
||
|
||
<!-- GENERATED: branches:start -->
|
||
<!-- GENERATED: branches:end -->
|
||
|
||
## 관련 일일 노트
|
||
|
||
해당 없음.
|
||
|
||
## 완료 후 정리
|
||
|
||
- PR 링크:
|
||
- 리뷰 메모:
|
||
- 머지 결과 / 배포 환경:
|
||
- **wiki 추출 대상** (verified만, `wiki/projects/`로만 추출):
|
||
- **추출하지 않을 항목** (planned / documented-only / abandoned):
|