Blog on Iterating

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

문제 상황

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

각 서비스는 각자의 NATS Account를 사용하고 있고, 자신의 메시지를 처리하는 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의 source/mirror 같은 기능을 사용해야 한다.

첫 번째 방법: Subject export/import - PUB 자체를 다른 Account로 전달하기

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

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

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

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

Account C에서는 이를 import한다.

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

그리고 Account C의 JetStream Stream이 해당 subject를 저장하게 만든다.

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

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

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

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

정상적인 상황에서는 Stream A와 Stream B에 쌓이는 메시지가 Stream C에도 쌓이는 것처럼 보인다.

문제는 저장 보장이다

Account A, B의 서비스는 Stream 저장을 위해 JetStream Publish를 사용해서 Publish를 할 것이고 PubAck를 이용할 것이다.

그런데 Account C의 경우에는 JetStream을 이용하는 것이 아니고 subject에 대한 구독을 sub 이용 하는 것이기 때문에 메시지가 스트림에 저장되는 것을 보장하지 못한다.

두 번째 방법: JetStream Mirror/Source

이 문제에 더 적합한 기능이 JetStream의 mirror/source다.
mirror에 대해서는 설명을 생략한다. mirror는 1:1 스트림 연결, source는 n:1 스트림 연결로 이해하면 된다.

export/import와 핵심 차이는 PUB을 복사하는 것이 아니라 이미 JetStream에 저장된 메시지를 복제한다는 것이다.

우리의 목표는 다음 구조를 만드는 것이다.

Service A               Service B
 │                       │
 ▼                       ▼
PUB - EVENTS_A          PUB - EVENTS_B
 │                       │ 
 ▼                       ▼
Stream A                Stream B           
 │                       │
 │   persisted message   │ 
 │                       │
 └───────────┬───────────┘                         
             ▼
           source
             │
             ▼
          Stream C

Source는 어떻게 동작하는가

source를 설정하면 복제 Stream이 원본 Stream을 단순 SUB하는 것이 아니다. 내부적으로 원본 Stream에 Consumer를 생성해 데이터를 전달 받는다.

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

같은 Account 안에 있는 Stream이라면 이 과정을 내부적으로 처리할 수 있다.

문제는 Stream A, Stream B, Stream C가 서로 다른 Account에 있다는 것이다. Account는 서로 격리되어 있기 때문에 Account C에서 Account A의 JetStream API를 호출하거나 source용 메시지를 전달받을 수 없다.

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

Source 구성

다른 계정간 Stream Source 설정을 위해서는 다음 세 종류의 통신 경로가 필요하다.

1. Consumer API - 원본 Stream에 source용 Consumer 생성 요청
2. Message Delivery - 실제 Stream 메시지 전달
3. Flow Control - 메시지 전달 속도 제어

Account A의 Stream_A를 Account C의 Stream에서 source한다고 가정하고 하나씩 구성해보자.

1. Consumer API를 Export 한다

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

여기서 먼저 짚어야 할 점은 stream exportstream은 JetStream의 Stream을 의미하지 않는다는 것이다. NATS Account의 export에는 크게 stream exportservice 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가 Account A의 특정 subject로 요청을 보내고 응답을 받을 수 있게 하는 방식입니다. 즉 request/reply 통신을 위한 export입니다.

Consumer Create API는 단순 메시지 전달이 아니라 요청에 대한 응답을 받는 JetStream API 호출이기 때문에 service export를 사용한다.

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

1번 Consumer API는 source용 Consumer를 생성하기 위한 Control Plane이다.

실제 Stream 데이터는 별도의 delivery subject를 통해 전달된다.

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

a.subject

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

a.subject.>

따라서 Account A에서는 해당 subject를 Account C에 export한다.

# Account A

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

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

# Account C

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

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

Account A                         Account C
Stream_A
   │
   ▼
Source Consumer
   │
   │ a.subject.*
   └────────────────────────────▶ Source ─────▶ Stream C

뒤에서 Source를 설정할 때 이 subject를 external.deliver에 지정한다.

external.deliver = a.subject

a.subject는 Account A의 Source Consumer가 Account C로 메시지를 전달할 때 사용할 delivery subject의 기준이 되는 값이라고 이해하면 된다.

3. Flow Control을 export한다

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

Push 기반으로 대량의 메시지가 전달될 경우 원본이 보내는 속도와 destination이 처리할 수 있는 속도에 차이가 발생할 수 있다.

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

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

$JS.FC.<STREAM>.>

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

$JS.FC.EVENTS_A.>

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

# Account A

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

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

Account A                       Account C

EVENTS_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는 자신의 JetStream API에서 사용하고 있으므로, Account A의 API는 다른 subject로 매핑해서 가져온다.

# 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. Stream C에 Source를 설정한다

지금까지 Account C에서 Source가 동작하기 위해 필요한 통신 경로를 연결했다.

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

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

이제 Account C의 Stream_C에서 Stream_A를 Source로 설정하면서 이 값들을 지정한다.

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

external.api에는 앞에서 Account C로 import한 Consumer API의 subject를 지정한다.

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

a.subject는 Account A에 생성되는 Source Consumer가 메시지를 전달할 때 사용하는 subject다. Account C에서 이 메시지를 받을 수 있도록 a.subject.>를 Account A에서 export하고 Account C에서 import한다.

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

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

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

Account B의 Stream_B도 같은 방식으로 설정하면 Stream_C에서 두 Stream을 함께 Source할 수 있다.

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

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

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