Blog on Iterating

NATS Account를 넘어서 JetStream 메시지를 안전하게 공유하기

문제 상황

현재 여러 서비스에서 NATS를 사용하고 있다.

각 서비스는 서로 다른 NATS Account를 사용하며, 각자의 Stream과 Consumer를 통해 메시지를 처리한다.

Service A ──> Stream A ──> Consumer A
Service B ──> Stream B ──> Consumer B

그런데 새로운 요구사항이 생겼다.

Service A와 Service B에서 발생하는 메시지를 하나의 Consumer에서 공통으로 처리해야 한다.

Service A ─┐
           ├──> Stream C ──> Consumer C
Service B ─┘

문제는 두 서비스가 서로 다른 NATS Account를 사용한다는 점이다.

NATS의 Account는 서로 독립된 namespace다. 별도의 설정이 없다면 Account A에서 publish한 메시지를 Account B에서는 볼 수 없다.

Account A

PUB events.created
Account B

SUB events.created   // 받을 수 없음

따라서 Account 간에 메시지를 전달하려면 export/import를 사용해야 한다. JetStream에 저장된 메시지를 다른 Account의 Stream으로 가져오려면 source/mirror도 함께 고려해야 한다.

첫 번째 방법: Subject export/import

가장 먼저 생각할 수 있는 방법은 메시지가 publish되는 subject 자체를 다른 Account에 공유하는 것이다.

NATS에서 export는 특정 subject에 대한 접근을 다른 Account에 공개하는 기능이고, import는 다른 Account가 공개한 subject를 자신의 Account에서 사용할 수 있도록 연결하는 기능이다.

Service A                Service B
   │                        │
   │ PUB a.events.>         │ PUB b.events.>
   ▼                        ▼
Account A                Account B
   │                        │
   ├── Stream A             ├── Stream B
   │                        │
   └── stream export        └── stream export 
           │                        │ 
           │                        │
           └───────────┬────────────┘
                       │             
                       ▼
                   Account C
                       │
                       ▼
                   Stream C

Account A에서는 다음과 같이 subject를 export한다.

exports: [
    { stream: "a.events.>", accounts: [C] }
]

Account C에서는 이를 import한다.

imports: [
    {
        stream: {
            account: A,
            subject: "a.events.>"
        }
    }
]

그리고 Account C의 JetStream Stream이 해당 subject를 저장하도록 설정한다.

{
  "name": "COMMON_EVENTS",
  "subjects": [
    "a.events.>"
  ]
}

이렇게 하면 새로 발생하는 메시지는 다음처럼 저장된다.

                ┌──> Stream A
PUB a.events ───┤
                └──> export/import ──> Stream C

                ┌──> Stream B
PUB b.events ───┤
                └──> export/import ──> Stream C

겉으로 보기에는 Stream A와 Stream B에 쌓이는 메시지가 Stream C에도 함께 저장되는 것처럼 보인다.

하지만 여기에는 중요한 차이가 있다. 저장을 보장하지 않는다는 것이다.

Account A와 B의 서비스는 JetStream Publish를 사용하고 PubAck를 받을 수 있다. 따라서 애플리케이션 입장에서는 메시지가 원본 Stream에 정상적으로 저장됐는지 확인할 수 있다.

반면 Account C로 전달되는 경로는 원본 Stream에 저장된 메시지를 복제하는 구조가 아니라, publish된 subject를 다른 Account에서도 subscribe할 수 있게 만든 구조다.

즉 원본 Stream의 저장 성공과 Account C의 Stream 저장 성공이 하나의 JetStream Publish ACK로 묶이지 않는다.

따라서 목표가 단순히 실시간 메시지를 다른 Account에서도 받는 것이 아니라, 원본 JetStream에 저장된 메시지를 기준으로 다른 Account에도 안정적으로 복제하는 것이라면 다른 방법이 필요하다.

두 번째 방법: JetStream Mirror/Source

이 요구사항에는 JetStream의 mirror/source가 더 잘 맞는다.

여기서는 source를 중심으로 살펴본다.

간단히 구분하면 다음과 같이 이해할 수 있다.

  • mirror: 하나의 Stream을 다른 Stream으로 1:1 복제

  • source: 여러 Stream을 하나의 Stream으로 가져올 수 있는 구조

Subject export/import와의 핵심적인 차이는 publish 자체를 공유하는 것이 아니라 이미 JetStream에 저장된 메시지를 복제한다는 점이다.

우리가 원하는 구조는 다음과 같다.

Service A               Service B
 │                       │
 ▼                       ▼
PUB - a.event           PUB - b.event
 │                       │ 
 ▼                       ▼
Stream A                Stream B           
 │                       │
 │   persisted message   │ 
 │                       │
 └───────────┬───────────┘                         
             ▼
           source
             │
             ▼
          Stream C

Source는 어떻게 동작하는가

source를 설정하면 복제 Stream이 원본 Stream을 단순 SUB하는 것이 아니다.

JetStream은 내부적으로 원본 Stream에 source용 Consumer를 생성하고, 이 Consumer를 통해 메시지를 전달받는다.

Stream A                                      Stream B  
┌────────────┐                                ┌────────────┐
│ messages A │                                │ messages B │
└────┬───────┘                                └────┬───────┘
     ▼                                             ▼
┌─────────────────┐                           ┌─────────────────┐
│ Consumer C Of A │                           │ Consumer C of B │
└────┬────────────┘                           └────┬────────────┘
     │ delivery                                    │ delivery
     └────────────────────┐   ┌────────────────────┘
                          ▼   ▼
                         Stream C   
                    ┌───────────────┐                                 
                    │ messages A, B │                                 
                    └───────────────┘
                     

모든 Stream이 같은 Account에 있다면 JetStream이 이 과정을 내부적으로 처리할 수 있다.

하지만 여기서는 Stream A, Stream B, Stream C가 서로 다른 Account에 있다.

NATS Account는 서로 격리되어 있으므로 Account C에서는 기본적으로 Account A의 JetStream API를 호출할 수도 없고, Account A의 source용 Consumer가 전달하는 메시지를 받을 수도 없다.

따라서 Cross Account Source를 사용하려면 source 동작에 필요한 subject들을 Account 사이에서 명시적으로 export/import해야 한다.

Cross Account Source 구성

Account A의 Stream_A를 Account C의 Stream_C에서 source한다고 가정해보자.

Cross Account Source가 동작하려면 크게 세 종류의 통신 경로가 필요하다.

  1. Consumer API - 원본 Stream에 source용 Consumer를 생성하기 위한 요청

  2. Message Delivery - 생성된 Consumer가 실제 Stream 메시지를 전달하는 경로

  3. Flow Control - 메시지 전달 속도를 제어하기 위한 통신 경로

하나씩 살펴보자.

1. Consumer 생성 요청 경로 열기

Account C에서 source를 생성하려면 원본인 Account A의 Stream_A에 source용 Consumer를 생성할 수 있어야 한다.

따라서 Account A에서는 Stream_A의 Consumer Create API를 export한다.

# Account A

exports: [
    {
        service: "$JS.API.CONSUMER.CREATE.Stream_A",
        accounts: [C]
    }
]

여기서 중요한 점은 stream export가 아니라 service export를 사용한다는 것이다. (위 exports 항목에 service 프로퍼티를 사용함)

stream export vs service export

먼저 혼동하기 쉬운 부분이 있다.

NATS Account의 stream export에서 말하는 stream은 JetStream의 Stream을 의미하지 않는다.

Account export에는 크게 stream export와 service export가 있으며, 둘은 통신 패턴에 따라 구분된다.

# stream export
Account A  ─────────────▶  Account C

           메시지 전달

# service export
Account C  ── request ──▶  Account A
           ◀─ response ──

           메시지 요청 + 응답

stream export는 Account A에서 publish되는 메시지를 다른 Account가 subscribe할 수 있도록 공유하는 방식이다.

즉, 일방향 pub/sub 통신이다.

반면 service export는 다른 Account가 특정 subject로 요청을 보내고 응답을 받을 수 있도록 한다.

즉, request/reply 통신을 위한 export다.

Consumer Create API는 단순한 메시지 전달이 아니라 JetStream API에 요청을 보내고 응답을 받아야 하므로 service export를 사용한다.

2. 메시지가 전달될 Subject를 Export한다

앞에서 구성한 Consumer API는 source용 Consumer를 만들기 위한 Control Plane이다.

실제 Stream의 데이터는 별도의 subject를 통해 전달된다. 이것을 delivery prefix라고 한다.

여기서 사용하는 delivery prefix는 애플리케이션이 publish하는 subject와는 별개다. JetStream Source가 내부적으로 메시지를 전달하기 위해 사용하는 경로다.

예를 들어 Account A에서 Account C로 source 메시지를 전달하기 위해 다음 delivery prefix를 사용한다고 하자.

a.delivery

Source Consumer가 실제 메시지를 전달할 때는 이 prefix 아래에 실제 Source Consumer의 delivery subject를 생성한다.

a.delivery.<generated>

따라서 Account A에서는 해당 delivery prefix의 하위 범위를 Account C에 export한다. 서비스에서 Account A로 pub할 때 사용하는 subject와 다른 것이다.

# Account A

{
    stream: "a.delivery.>",
    accounts: [C]
}

그리고 Account C에서는 이를 import한다.

# Account C

{
    stream: {
        account: A,
        subject: "a.delivery.>"
    }
}

이제 Account A의 Source Consumer가 전달하는 메시지가 Account C까지 이동할 수 있다.

Account A                         Account C
Stream_A
   │
   ▼
Source Consumer
   │
   │ a.delivery.>
   └────────────────────────────▶ Source ─────▶ Stream C

뒤에서 Source를 설정할 때 이 a.delivery 값을 external.deliver에 지정한다.

external.deliver = a.delivery

즉 a.delivery는 Account A에 만들어진 Source Consumer가 Account C로 메시지를 전달할 때 사용하는 delivery prefix라고 이해하면 된다.

3. Flow Control을 Export한다

Source는 내부적으로 push consumer를 사용해 메시지를 전달받는다.

Push 방식으로 대량의 메시지를 전송하다 보면 원본이 메시지를 보내는 속도와 destination이 처리하는 속도 사이에 차이가 발생할 수 있다.

JetStream은 이를 조절하기 위해 Flow Control을 사용한다.

Stream의 Flow Control subject는 다음과 같은 형태다.

$JS.FC.<STREAM>.>

따라서 Stream_A의 Flow Control subject는 다음과 같다.

$JS.FC.Stream_A.>

이 subject 역시 Account A와 Account C 사이에서 통신할 수 있도록 열어줘야 한다.

# Account A

{
    service: "$JS.FC.Stream_A.>",
    accounts: [C]
}

Flow Control 역시 Account C에서 Account A 방향으로 요청을 전달해야 하므로 Account A의 service export뿐 아니라 Account C의 service import도 필요하다.

# Account C

{
    service: {
        account: A,
        subject: "$JS.FC.Stream_A.>"
    }
}

Flow Control은 Source가 지속적으로 메시지를 전달받기 위한 제어 경로라고 볼 수 있다.

Account A                       Account C
Stream_A                         Source
   │                               │
   │      Message Delivery         │
   ├──────────────────────────────>│
   │                               │
   │       Flow Control            │
   │<─────────────────────────────>│

4. Account C에서 JetStream API를 Import한다

앞에서 Account A는 Source Consumer를 생성할 수 있도록 다음 Consumer API를 export했다.

$JS.API.CONSUMER.CREATE.Stream_A

이제 Account C에서 이 API를 import해야 한다.

그런데 Account C 역시 $JS.API namespace를 자신의 JetStream API에서 사용하고 있다.

따라서 Account A의 JetStream API에 다음의 별칭을 부여한다.

$JS.A.API.CONSUMER.CREATE.Stream_A

Account C에 import 방식은 다음과 같다.

# Account C

{
    service: {
        account: A,
        subject: "$JS.API.CONSUMER.CREATE.Stream_A"
    },
    to: "$JS.A.API.CONSUMER.CREATE.Stream_A"
}

이렇게 하면 Account C에서

$JS.A.API.CONSUMER.CREATE.Stream_A

로 요청한 메시지는 Account A의

$JS.API.CONSUMER.CREATE.Stream_A

로 전달된다.

5. Source 설정과 연결한다

지금까지 Account C에서 Cross Account Source가 동작하는 데 필요한 통신 경로를 구성했다.

Source 설정에서 직접 지정하는 값은 다음 두 가지다.

Consumer API      → $JS.A.API
Message Delivery  → a.delivery

이제 Account C의 Stream_C에서 Stream_A를 Source로 설정한다.

# Account C
{
  "name": "Stream_C",
  "sources": [
    {
      "name": "Stream_A",
      "external": {
        "api": "$JS.A.API",
        "deliver": "a.delivery"
      }
    }
  ]
}

external.api에는 Account C로 import한 JetStream API의 prefix를 지정한다.

Source는 이 경로를 통해 Account A의 JetStream API에 접근하고 Stream_A에 source용 Consumer를 생성한다.

external.deliver에는 Source Consumer의 실제 delivery subject를 생성할 때 사용할 delivery prefix를 지정한다.

앞서 a.delivery.>를 Account A에서 export하고 Account C에서 import했기 때문에 Source Consumer가 전달하는 메시지는 Account C까지 이동할 수 있다.

Account A                              Account C
Stream_A
   │
   ▼
Source Consumer
   │
   │ a.delivery.>
   └─────────────────────────────────▶ Source
                                         │
                                         ▼
                                      Stream_C

정리하면 다음과 같다.

external.api
    → Source Consumer를 생성하기 위한 API 경로

external.deliver
    → 생성된 Consumer가 실제 메시지를 전달하기 위한 경로

6. 여러 Account의 Stream을 하나로 모으기

Account B의 Stream_B도 같은 방식으로 구성하면 Account C의 Stream_C에서 두 Stream을 함께 source할 수 있다.

{
  "name": "Stream_C",
  "sources": [
    {
      "name": "Stream_A",
      "external": {
        "api": "$JS.A.API",
        "deliver": "a.delivery"
      }
    },
    {
      "name": "Stream_B",
      "external": {
        "api": "$JS.B.API",
        "deliver": "b.delivery"
      }
    }
  ]
}

최종적인 구조는 다음과 같다.

Account A                              Account B

Stream_A                               Stream_B
   │                                      │
   ▼                                      ▼
Source Consumer                       Source Consumer
   │                                      │
   │ a.delivery.>                         │ b.delivery.>
   └────────────────┐    ┌────────────────┘
                    ▼    ▼
                   Account C
                       │
                       ▼
                    Stream_C
                       │
                       ▼
                   Consumer C

Account C의 Stream_C는 Account A와 Account B의 Stream에 이미 저장된 메시지를 각각 Source Consumer를 통해 전달받아 하나의 Stream에 저장한다.

정리

Cross Account 환경에서 여러 JetStream 메시지를 하나의 Stream으로 모으려면 단순한 subject export/import와 JetStream Source의 차이를 이해하는 것이 중요하다.

Subject export/import는 publish되는 메시지의 subject를 다른 Account에 공개하는 방식이다.

반면 JetStream Source는 원본 Stream에 저장된 메시지를 Consumer를 통해 destination Stream으로 복제하는 방식이다.

그리고 Source 대상이 다른 Account에 있다면 Account 격리 때문에 몇 가지 추가 경로를 직접 연결해야 한다.

Consumer API
    → source용 Consumer 생성

Message Delivery
    → 실제 메시지 전달

Flow Control
    → 전달 속도 제어

결국 Cross Account Source를 구성한다는 것은 단순히 Stream 하나를 다른 Account에 공개하는 작업이 아니다.

JetStream Source가 내부적으로 Consumer를 생성하고 메시지를 전달하는 데 필요한 Control Plane과 Data Plane을 Account 사이에 연결하는 작업이라고 이해하면 전체 구조가 훨씬 명확해진다.