콘텐츠로 이동

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

중앙 BOM 보드에서 호환성, 데이터, 메시징, 운영 안전성 작업대로 경로가 갈라지는 3D 미니어처 작업실
전달, 업무 성공, acknowledgement, checkpoint는 서로 다른 사건이며 소유자도 다릅니다.

메시지가 핸들러에 도착했다고 업무 처리가 끝난 것은 아닙니다. 데이터베이스에 쓴 뒤 ack 전에 프로세스가 종료될 수 있고, 스트림 레코드를 방출한 뒤 체크포인트 저장이 실패할 수도 있습니다. 재전달 가능한 시스템에서는 이 틈을 정상 경로로 설계해야 합니다.

bluetape4k-dependencies 2.0.0bluetape4k-projects 2.0.0bluetape4k-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를 하나의 원자적 트랜잭션으로 만들어 주지 않습니다.

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과 체크포인트를 중단합니다.

AdapterOrderingProgress state호출자 책임
DynamoDB Streamsparent-before-childinclusive checkpoint replay멱등성, 체크포인트 저장소 수명
Kinesisshard graph와 parent orderinglease·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와 자동으로 상호운용되지 않습니다.

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 issueevent externalization visual companion이 경계를 보여 줍니다.

대상어댑터 소유호출자 소유
NATS Flowcold 브리지, 용량 제한 채널, 생성한 handle 정리연결, 업무 성공, ack 정책
AWS Streamsshard 탐색, 방출 순서, 정의된 체크포인트 hook멱등성, 영속 저장소, 운영 설정
SQS listenerpolling·heartbeat·partial ack 프로토콜핸들러 트랜잭션, retry/DLQ, 페이로드 저장 정책
SNS HTTP허용 목록·서명 검증허용 topic, confirmation, 업무 부수 효과
Modulith선택한 송신/수신 어댑터이벤트 계약, 버전, 멱등성, 반대 방향 활성화

어댑터가 만든 handle만 어댑터가 닫습니다. 클라이언트, 자격 증명, IAM, queue/topic, 영속 저장소, 업무 트랜잭션은 애플리케이션 경계에 남습니다.

댓글

GitHub 계정으로 의견을 남기거나 reaction을 남길 수 있습니다.