NATS 테이블 엔진

NATS 테이블 엔진

이 엔진은 ClickHouse를 NATS와 통합할 수 있게 해줘요. NATS를 사용하면 다음을 할 수 있어요.

  • 메시지 주제(subject)를 게시하거나 구독하기
  • 새 메시지가 생기면 바로 처리하기

출처: 문서

본문

테이블 생성하기

CREATE TABLE [IF NOT EXISTS] [db.]table_name [ON CLUSTER cluster]
(
    name1 [type1] [DEFAULT|MATERIALIZED|ALIAS expr1],
    name2 [type2] [DEFAULT|MATERIALIZED|ALIAS expr2],
    ...
) ENGINE = NATS SETTINGS
    nats_url = 'host:port',
    nats_subjects = 'subject1,subject2,...',
    nats_format = 'data_format'[,]
    [nats_schema = '',]
    [nats_num_consumers = N,]
    [nats_queue_group = 'group_name',]
    [nats_secure = false,]
    [nats_max_reconnect = N,]
    [nats_reconnect_wait = N,]
    [nats_server_list = 'host1:port1,host2:port2,...',]
    [nats_skip_broken_messages = N,]
    [nats_max_block_size = N,]
    [nats_flush_interval_ms = N,]
    [nats_username = 'user',]
    [nats_password = 'password',]
    [nats_token = 'clickhouse',]
    [nats_credentials = '-----BEGIN NATS USER JWT----- ...',]
    [nats_startup_connect_tries = 5,]
    [nats_max_rows_per_message = 1,]
    [nats_commit_on_select = false,]
    [nats_handle_error_mode = 'default']

필수 매개변수:

  • nats_url – host:port (예: localhost:4222)
  • nats_subjects – NATS 테이블이 구독/게시할 주제 목록. foo.*.bar 또는 baz.> 같은 와일드카드 주제를 지원해요
  • nats_format – 메시지 형식. SQL FORMAT 함수와 같은 표기법을 사용해요 (예: JSONEachRow). 자세한 내용은 Formats 섹션을 참고해요

선택 매개변수:

  • nats_schema – 형식이 스키마 정의를 요구할 때 반드시 사용해야 하는 매개변수. 예를 들어 Cap'n Proto는 스키마 파일의 경로와 루트 schema.capnp:Message 객체의 이름을 요구해요
  • nats_stream – NATS JetStream에 존재하는 스트림의 이름
  • nats_consumer_name – NATS JetStream에 존재하는 durable pull 소비자의 이름
  • nats_num_consumers – 테이블당 소비자 수. 기본값: 1. 한 소비자의 처리량이 부족하면 NATS core에 대해 더 많은 소비자를 지정해요
  • nats_queue_group – NATS 구독자의 큐 그룹 이름. 기본값은 테이블 이름
  • nats_max_reconnect – 더 이상 사용되지 않으며 효과가 없어요. 재연결은 nats_reconnect_wait 타임아웃으로 영구적으로 수행돼요
  • nats_reconnect_wait – 각 재연결 시도 사이에 대기하는 밀리초 단위 시간. 기본값: 2000
  • nats_server_list – 연결을 위한 서버 목록. NATS 클러스터에 연결하려면 지정할 수 있어요
  • nats_skip_broken_messages – 블록당 스키마와 맞지 않는 메시지에 대한 NATS 메시지 파서 허용치. 기본값: 0. nats_skip_broken_messages = N이면 엔진은 파싱할 수 없는 NATS 메시지 N개를 건너뛰어요 (메시지 하나는 데이터 행 하나와 같음)
  • nats_max_block_size – NATS에서 데이터를 flush하기 위해 poll이 모은 행 수. 기본값: max_insert_block_size
  • nats_flush_interval_ms – NATS에서 읽은 데이터를 flush하는 타임아웃. 기본값: stream_flush_interval_ms
  • nats_wait_for_flush_intervaltrue면 백그라운드 스트리밍 주기가 소비자 큐가 비워지는 즉시 끝나는 대신 전체 flush 간격(nats_flush_interval_ms, 그렇지 않으면 stream_flush_interval_ms) 동안 열려 있어 더 많은 메시지를 단일 블록에 축적할 수 있고, 대가로 최대 한 번의 flush 간격의 추가 수집 지연이 발생해요. 기본값: false (저지연 drain-and-go 동작)
  • nats_username – NATS 사용자 이름. 서버 구성 파일에 정의된 named collection에 저장되면 쿼리가 컬렉션의 nats_url이나 nats_server_list를 재정의할 수 없어요
  • nats_password – NATS 비밀번호. 서버 구성 파일에 정의된 named collection에 저장되면 쿼리가 컬렉션의 nats_url이나 nats_server_list를 재정의할 수 없어요
  • nats_token – NATS 인증 토큰. 서버 구성 파일에 정의된 named collection에 저장되면 쿼리가 컬렉션의 nats_url이나 nats_server_list를 재정의할 수 없어요
  • nats_credential_file – NATS 자격 증명 파일 경로. 서버가 자신의 권한으로 경로를 열기 때문에 서버 구성 파일에 정의되고 nats_urlnats_server_list가 쿼리에 의해 재정의되지 않은 named collection에서만 허용돼요. 쿼리에서는 대신 nats_credentials에 파일 내용을 전달해요
  • nats_credentials – NATS 자격 증명 내용(사용자 JWT와 시드가 있는 .creds 파일과 같은 페이로드). 쿼리가 사용할 수 있는 유일한 표기법이므로, 운영자가 <nats_credential_file overridable="false">로 해당 경로를 잠그지 않는 한 named collection에서 상속된 nats_credential_file과 충돌하지 않고 대체해요. named collection이 담고 있는 자격 증명을 버리기 위해 빈 문자열을 할당할 수는 없어요
  • nats_ca_file – NATS 서버 인증서를 검증하는 데 사용되는 신뢰된 CA 인증서가 있는 파일 경로. nats_secure가 필요해요. nats_credential_file처럼 서버가 자신의 권한으로 경로를 열기 때문에 서버 구성 파일에 정의되고 nats_urlnats_server_list가 쿼리에 의해 재정의되지 않은 named collection에서만 허용돼요
  • nats_client_cert_file – NATS 서버에 제시하는 클라이언트 인증서 경로. nats_securenats_client_key_file이 필요해요. nats_ca_file과 같은 출처에서 허용돼요
  • nats_client_key_filenats_client_cert_file의 개인 키 경로. nats_ca_file과 같은 출처에서 허용돼요
  • nats_startup_connect_tries – 시작 시 연결 시도 횟수. 기본값: 5
  • nats_max_rows_per_message — 행 기반 형식에서 하나의 NATS 메시지에 쓰는 최대 행 수. (기본값: 1)
  • nats_commit_on_select – 쿼리가 만들어질 때 메시지를 커밋해요. JetStream에만 적용되고, core NATS는 확인(acknowledgement)이 없어요. 기본값: 0
  • nats_handle_error_mode — NATS 엔진의 오류 처리 방식. 가능한 값: default(메시지 파싱 실패 시 예외 발생), stream(예외 메시지와 원시 메시지가 가상 컬럼 _error_raw_message에 저장)

SSL 연결:

보안 연결에는 nats_secure = 1을 사용해요. 인증서 검증은 CLICKHOUSE_NATS_TLS_SECURE 환경 변수로 제어돼요. 인증서가 만료, 자체 서명, 누락이거나 다른 방식으로 무효하면 CLICKHOUSE_NATS_TLS_SECURE=0으로 설정해 검증을 비활성화해요. 개인 CA가 서명한 서버 인증서는 검증을 끄는 것보다 nats_ca_file을 CA 인증서로 지정해 검증하는 것이 더 좋아요. 서버가 클라이언트 인증서를 요구하면 nats_client_cert_filenats_client_key_file로 제공해요. 세 가지 모두 운영자 설정이에요: 서버 구성 파일에 정의된 named collection에서 옵니다. 테이블이 연결할 때마다 각 파일이 읽히므로, 읽을 수 없거나 형식이 잘못된 파일이 있으면 핸드셰이크 대신 쿼리를 실패시켜요.

NATS 테이블에 쓰기:

테이블이 단일 주제에서만 읽으면 어떤 삽입도 같은 주제에 게시해요. 그러나 테이블이 여러 주제에서 읽으면 어떤 주제에 게시할지 지정해야 해요. 그렇기 때문에 여러 주제가 있는 테이블에 삽입할 때는 stream_like_engine_insert_queue 설정이 필요해요. 테이블이 읽는 주제 중 하나를 선택해 그곳에 데이터를 게시할 수 있어요. 예:

CREATE TABLE queue (
    key UInt64,
    value UInt64
  ) ENGINE = NATS
    SETTINGS nats_url = 'localhost:4444',
             nats_subjects = 'subject1,subject2',
             nats_format = 'JSONEachRow';

INSERT INTO queue
SETTINGS stream_like_engine_insert_queue = 'subject2'
VALUES (1, 1);

또한 nats 관련 설정과 함께 형식 설정을 추가할 수 있어요.

예시:

CREATE TABLE queue (
    key UInt64,
    value UInt64,
    date DateTime
  ) ENGINE = NATS
    SETTINGS nats_url = 'localhost:4444',
             nats_subjects = 'subject1',
             nats_format = 'JSONEachRow',
             date_time_input_format = 'best_effort';

NATS 서버 구성은 ClickHouse 구성 파일로 추가할 수 있어요. 더 구체적으로 NATS 엔진의 비밀번호를 추가할 수 있어요.

<nats>
    <user>click</user>
    <password>house</password>
    <token>clickhouse</token>
</nats>

설명

SELECT는 (디버깅 외에는) 메시지 읽기에 특히 유용하지 않아요. 각 메시지는 한 번만 읽을 수 있기 때문이에요. materialized view를 사용해 실시간 스레드를 만드는 것이 더 실용적이에요. 그러려면:

  • 엔진을 사용해 NATS 소비자를 만들고 그것을 데이터 스트림으로 간주해요
  • 원하는 구조로 테이블을 만들어요
  • 엔진의 데이터를 변환해 앞서 만든 테이블에 넣는 materialized view를 만들어요

MATERIALIZED VIEW가 엔진에 조인되면 백그라운드에서 데이터 수집을 시작해요. 이렇게 하면 NATS에서 메시지를 계속 받아 SELECT로 필요한 형식으로 변환할 수 있어요. 하나의 NATS 테이블은 원하는 만큼 많은 materialized view를 가질 수 있어요. view는 테이블에서 직접 데이터를 읽지 않고 새 레코드(블록 단위)를 받으므로, 여러 테이블에 다른 세부 수준으로 쓸 수 있어요 (그룹화-집계 포함, 포함하지 않고).

예시:

CREATE TABLE queue (
    key UInt64,
    value UInt64
  ) ENGINE = NATS
    SETTINGS nats_url = 'localhost:4444',
             nats_subjects = 'subject1',
             nats_format = 'JSONEachRow',
             date_time_input_format = 'best_effort';

CREATE TABLE daily (key UInt64, value UInt64)
    ENGINE = MergeTree() ORDER BY key;

CREATE MATERIALIZED VIEW consumer TO daily
    AS SELECT key, value FROM queue;

SELECT key, value FROM daily ORDER BY key;

스트림 데이터 수신을 중지하거나 변환 로직을 바꾸려면 materialized view를 분리해요:

DETACH TABLE consumer;
ATTACH TABLE consumer;

ALTER로 대상 테이블을 바꾸고 싶다면 대상 테이블과 view의 데이터 사이 괴리를 피하기 위해 materialized view를 비활성화하는 것을 권장해요.

가상 컬럼

  • _subject - NATS 메시지 주제. 데이터 타입: String

nats_handle_error_mode='stream'일 때 추가 가상 컬럼:

  • _raw_message - 성공적으로 파싱되지 못한 원시 메시지. 데이터 타입: Nullable(String)
  • _error - 실패한 파싱 중 발생한 예외 메시지. 데이터 타입: Nullable(String)

참고: _raw_message_error 가상 컬럼은 파싱 중 예외가 발생한 경우에만 채워지며, 메시지가 성공적으로 파싱되면 항상 NULL이에요.

데이터 형식 지원

NATS 엔진은 ClickHouse에서 지원하는 모든 형식을 지원해요. 하나의 NATS 메시지의 행 수는 형식이 행 기반인지 블록 기반인지에 따라 달라져요.

  • 행 기반 형식의 경우 하나의 NATS 메시지의 행 수를 nats_max_rows_per_message 설정으로 제어할 수 있어요
  • 블록 기반 형식의 경우 블록을 더 작은 부분으로 나눌 수 없지만, 하나의 블록의 행 수는 일반 설정 max_block_size로 제어할 수 있어요

JetStream 사용하기

NATS 엔진을 NATS JetStream과 사용하기 전에 NATS 스트림과 durable pull 소비자를 만들어야 해요. 이를 위해 예를 들어 NATS CLI 패키지의 nats 유틸리티를 사용할 수 있어요.

스트림 생성 `$ nats stream add ? Stream Name stream_name ? Subjects stream_subject ? Storage file ? Replication 1 ? Retention Policy Limits ? Discard Policy Old ? Stream Messages Limit -1 ? Per Subject Messages Limit -1 ? Total Stream Size -1 ? Message TTL -1 ? Max Message Size -1 ? Duplicate tracking time window 2m0s ? Allow message Roll-ups No ? Allow message deletion Yes ? Allow purging subjects or the entire stream Yes Stream stream_name was created

Information for Stream stream_name created 2025-10-03 14:12:51

            Subjects: stream_subject
            Replicas: 1
             Storage: File

Options:

           Retention: Limits
     Acknowledgments: true
      Discard Policy: Old
    Duplicate Window: 2m0s
          Direct Get: true
   Allows Msg Delete: true
        Allows Purge: true

Allows Per-Message TTL: false Allows Rollups: false

Limits:

    Maximum Messages: unlimited
 Maximum Per Subject: unlimited
       Maximum Bytes: unlimited
         Maximum Age: unlimited
Maximum Message Size: unlimited
   Maximum Consumers: unlimited

State:

            Messages: 0
               Bytes: 0 B
      First Sequence: 0
       Last Sequence: 0
    Active Consumers: 0

`

durable pull 소비자 생성 `$ nats consumer add ? Select a Stream stream_name ? Consumer name consumer_name ? Delivery target (empty for Pull Consumers) ? Start policy (all, new, last, subject, 1h, msg sequence) all ? Acknowledgment policy explicit ? Replay policy instant ? Filter Stream by subjects (blank for all) ? Maximum Allowed Deliveries -1 ? Maximum Acknowledgments Pending 0 ? Deliver headers only without bodies No ? Add a Retry Backoff Policy No Information for Consumer stream_name > consumer_name created 2025-10-03T14:13:51+03:00

Configuration:

                Name: consumer_name
           Pull Mode: true
      Deliver Policy: All
          Ack Policy: Explicit
            Ack Wait: 30.00s
       Replay Policy: Instant
     Max Ack Pending: 1,000
   Max Waiting Pulls: 512

State:

Last Delivered Message: Consumer sequence: 0 Stream sequence: 0 Acknowledgment Floor: Consumer sequence: 0 Stream sequence: 0 Outstanding Acks: 0 out of maximum 1,000 Redelivered Messages: 0 Unprocessed Messages: 0 Waiting Pulls: 0 of maximum 512 `

스트림과 durable pull 소비자를 만든 뒤 NATS 엔진으로 테이블을 만들 수 있어요. 그러려면 nats_stream, nats_consumer_name, nats_subjects를 초기화해야 해요.

CREATE TABLE nats_jet_stream (
    key UInt64,
    value UInt64
  ) ENGINE NATS
    SETTINGS  nats_url = 'localhost:4222',
              nats_stream = 'stream_name',
              nats_consumer_name = 'consumer_name',
              nats_subjects = 'stream_subject',
              nats_format = 'JSONEachRow';

JetStream 테이블은 최소 한 번(at-least-once) 전달을 제공해요. 메시지는 의존하는 materialized view에 삽입된 후에만 확인(acknowledge)되므로, 삽입이 실패하거나 중단된 메시지는 확인되지 않은 상태로 남아 재전달돼요. core NATS(JetStream 없음)에는 확인이나 재생(replay)이 없으므로 최대 한 번(at-most-once)이며 중단된 메시지는 손실돼요.

데이터 내구성(Durability)

이 섹션은 JetStream에만 적용돼요. core NATS는 위에서 설명한 대로 확인이 없고 최대 한 번이므로, 확인된 메시지가 손실될 수 있는 창이 없어요.

JetStream 테이블은 삽입된 데이터가 디스크에 쓰이기 전에 OS 페이지 캐시가 버려지면 이미 소비된 행을 조용히 잃을 수 있어요. 배치가 의존하는 materialized view에 푸시된 후 소비자가 그 메시지를 확인(acknowledge)해 스트림이 그 메시지들을 넘어 진행돼요. 그러나 삽입된 행은 대상 파트가 fsync된 후에만 내구성이 생기는데, 기본적으로 동기적으로 발생하지 않아요 (fsync_after_insert = 0). 확인 후 대상 파트가 fsync되기 전에 페이지 캐시가 손실되면 메시지가 더 이상 재전달되지 않으므로 행이 오류 없이 손실되고 count()가 단순히 더 작아져요. 일반적인 프로세스 종료는 커널이 페이지 캐시를 유지하고 결국 기록하므로 이를 드러내지 않아요. 페이지 캐시 손실은 드러내요. 예로 장치 수준 전원 손실과 비정상적인 호스트나 커널 리셋이 있어요.

권장하는 materialized-view 소비 경로(전체 삽입 파이프라인이 끝난 뒤에만 확인을 보내는)에서는 대상 MergeTree 테이블에 fsync_after_insert = 1(그리고 fsync_part_directory = 1)을 설정하면 삽입된 파트가 확인 전에 내구성이 생겨 이 창이 크게 좁아져요. 이 설정은 배치가 삽입되는 모든 MergeTree 테이블, 캐스케이드된 materialized-view 대상까지 포함해 활성화해야 해요. 기본값으로 남겨둔 테이블은 여전히 파트를 잃을 수 있어요. 비동기 중개자는 이 설정만으로 내구성을 얻지 못해요. 예를 들어 distributed_foreground_insert = 0이면 Distributed 대상이 백그라운드에서 삽입하는데, 이는 ClickHouse Cloud 외부의 기본값이므로 자체 내구성 설정이나 동기 삽입이 필요해요. 이 완화는 nats_commit_on_select = 1을 사용한 직접 INSERT ... SELECT ... FROM <nats_table>에도 적용되지 않아요. 그 경우 메시지는 대상이 내구성 있는 파트를 쓴 뒤가 아니라 읽기가 끝에 도달할 때 확인되기 때문이에요.

더 알아보기 (Learn more)