bluetape4k-dependencies 2.0.0 활용기 Part 3: 메시징과 스트림의 소유권

메시지가 핸들러에 도착했다고 업무 처리가 끝난 것은 아닙니다. 데이터베이스에 쓴 뒤 ack 전에 프로세스가 종료될 수 있고, 스트림 레코드를 방출한 뒤 체크포인트 저장이 실패할 수도 있습니다. 재전달 가능한 시스템에서는 이 틈을 정상 경로로 설계해야 합니다.
bluetape4k-dependencies 2.0.0은
bluetape4k-projects 2.0.0과
bluetape4k-aws 1.0.0을 함께 선택합니다.
두 릴리스의 메시징 API는 at-least-once 재전달과 호출자가 책임지는 멱등성을 숨기지 않습니다.
메시지 수신과 업무 성공은 같은 사건이 아니다
섹션 제목: “메시지 수신과 업무 성공은 같은 사건이 아니다”가장 작은 안전 순서는 다음과 같습니다.
flow.collect { message -> try { processIdempotently(message) // 업무 처리의 부수 효과가 먼저 성공한다. message.ack() // 그다음 소비 위치를 전진시킨다. } catch (retryable: RetryableFailure) { message.nak() } catch (permanent: PermanentFailure) { message.term() }}실제 트랜잭션과 브로커 기능에 맞춰 inbox/outbox 또는 멱등성 키가 필요할 수 있습니다. 어댑터는 업무 처리의 부수 효과와 브로커 ack를 하나의 원자적 트랜잭션으로 만들어 주지 않습니다.
NATS 수동 ack
섹션 제목: “NATS 수동 ack”Projects 2.0.0은 JetStream ConsumerContext pull consumer와 JetStream push subscription을 수집할 때마다 생성하는 cold
Flow<Message>를 추가했습니다. Flow 용량과 NATS pending limit을 검증하며, 취소 시 어댑터가 만든
JetStreamSubscription 또는 IterableConsumer만 닫습니다. 연결과 stream/consumer 설정은 호출자가 소유합니다.
어댑터는 수동 acknowledgement만 제공합니다. 업무 성공 뒤 호출자가 ack()를 호출하고, 재시도할 실패에는 nak(),
처리할 수 없는 메시지에는 term() 정책을 적용합니다. 대기 큐에서 메시지가 누락되면 정상 완료가 아니라
NatsConsumerFlowException입니다.
DynamoDB Streams와 Kinesis 체크포인트
섹션 제목: “DynamoDB Streams와 Kinesis 체크포인트”DynamoDB Streams Flow는 shard graph의 parent-before-child 순서를 따릅니다. 체크포인트에서 다시 시작할 때는 저장된
sequence를 포함해 replay할 수 있으므로 consumer의 업무 처리는 멱등성을 보장해야 합니다. 이 어댑터에는 Kinesis와 같은
분산 lease 계약이 없습니다.
Kinesis consumerFlow는 shard discovery, parent ordering, bounded concurrency, lease와 fencing을 함께 다룹니다.
레코드를 downstream으로 emit한 것만으로 체크포인트를 저장하지 않습니다. 수집기가 해당 레코드를 처리하고 반환한 뒤에야
체크포인트를 전진시킵니다. Lease를 잃으면 새 emit과 체크포인트를 중단합니다.
| Adapter | Ordering | Progress state | 호출자 책임 |
|---|---|---|---|
| DynamoDB Streams | parent-before-child | inclusive checkpoint replay | 멱등성, 체크포인트 저장소 수명 |
| Kinesis | shard graph와 parent ordering | lease·fencing 뒤 체크포인트 | 멱등성, lease/checkpoint 저장소 구성 |
SQS의 visibility와 페이로드 소유권
섹션 제목: “SQS의 visibility와 페이로드 소유권”긴 작업은 visibility timeout보다 오래 걸릴 수 있습니다. SQS visibility heartbeat는 처리 중인 메시지의 visibility를 연장하지만 업무 성공을 대신 판정하지 않습니다. 배치 listener는 성공한 항목만 부분 acknowledgement하고 실패한 항목을 재전달 대상으로 남깁니다. FIFO queue에서는 메시지 그룹 순서와 backpressure도 함께 지켜야 합니다.
Extended Client는 큰 페이로드를 S3에 저장하고 SQS 본문에는 서명된 포인터를 둡니다. 애플리케이션은 bucket, IAM, encryption, retention, lifecycle, 재시도와 고아 객체 정리를 구성합니다. 이 포인터는 모듈 전용이며 AWS Java Extended Client와 자동으로 상호운용되지 않습니다.
- Visibility heartbeat
- Partial acknowledgement
- Extended Client payload offload
- SQS Extended Client visual companion
SNS 검증과 Modulith 방향
섹션 제목: “SNS 검증과 Modulith 방향”SNS HTTP 어댑터는 구조를 파싱한 뒤 정확히 일치하는 TopicArn 허용 목록을 네트워크 접근 전에 적용하고, 허용된 인증서 URL에서
인증서를 가져와 Signature v1/v2를 검증합니다. 검증을 통과한 메시지만 핸들러로 보냅니다. Subscription confirmation은
자동 부수 효과가 아니라 핸들러가 NotificationStatus.confirmSubscription()을 명시적으로 호출할 때 수행됩니다.
SNS signature visual companion에서 이 순서를 단계별로 확인할 수 있습니다.
Spring Modulith 연동은 두 방향을 구분합니다.
- Producer externalization: 애플리케이션 이벤트를 SNS/SQS로 내보냅니다.
- Consumer ingress: DIRECT body 또는 SNS-wrapped SQS body를 검증하고 애플리케이션 dispatch로 들여옵니다.
한 방향을 활성화했다고 반대 방향까지 자동 구성되는 것은 아닙니다. Modulith issue와 event externalization visual companion이 경계를 보여 줍니다.
소유권 표
섹션 제목: “소유권 표”| 대상 | 어댑터 소유 | 호출자 소유 |
|---|---|---|
| NATS Flow | cold 브리지, 용량 제한 채널, 생성한 handle 정리 | 연결, 업무 성공, ack 정책 |
| AWS Streams | shard 탐색, 방출 순서, 정의된 체크포인트 hook | 멱등성, 영속 저장소, 운영 설정 |
| SQS listener | polling·heartbeat·partial ack 프로토콜 | 핸들러 트랜잭션, retry/DLQ, 페이로드 저장 정책 |
| SNS HTTP | 허용 목록·서명 검증 | 허용 topic, confirmation, 업무 부수 효과 |
| Modulith | 선택한 송신/수신 어댑터 | 이벤트 계약, 버전, 멱등성, 반대 방향 활성화 |
어댑터가 만든 handle만 어댑터가 닫습니다. 클라이언트, 자격 증명, IAM, queue/topic, 영속 저장소, 업무 트랜잭션은 애플리케이션 경계에 남습니다.
bluetape4k-projects 2.0.0bluetape4k-aws 1.0.0- AWS messaging·security Epic
- SNS signature verification
댓글
GitHub 계정으로 의견을 남기거나 reaction을 남길 수 있습니다.