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– 메시지 형식. SQLFORMAT함수와 같은 표기법을 사용해요 (예: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– 각 재연결 시도 사이에 대기하는 밀리초 단위 시간. 기본값:2000nats_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_sizenats_flush_interval_ms– NATS에서 읽은 데이터를 flush하는 타임아웃. 기본값: stream_flush_interval_msnats_wait_for_flush_interval–true면 백그라운드 스트리밍 주기가 소비자 큐가 비워지는 즉시 끝나는 대신 전체 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_url과nats_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_url과nats_server_list가 쿼리에 의해 재정의되지 않은 named collection에서만 허용돼요nats_client_cert_file– NATS 서버에 제시하는 클라이언트 인증서 경로.nats_secure와nats_client_key_file이 필요해요.nats_ca_file과 같은 출처에서 허용돼요nats_client_key_file–nats_client_cert_file의 개인 키 경로.nats_ca_file과 같은 출처에서 허용돼요nats_startup_connect_tries– 시작 시 연결 시도 횟수. 기본값:5nats_max_rows_per_message— 행 기반 형식에서 하나의 NATS 메시지에 쓰는 최대 행 수. (기본값:1)nats_commit_on_select– 쿼리가 만들어질 때 메시지를 커밋해요. JetStream에만 적용되고, core NATS는 확인(acknowledgement)이 없어요. 기본값:0nats_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_file과 nats_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>에도 적용되지 않아요. 그 경우 메시지는 대상이 내구성 있는 파트를 쓴 뒤가 아니라 읽기가 끝에 도달할 때 확인되기 때문이에요.