Skip to main content
이 엔진을 사용하면 ClickHouse를 NATS와 통합할 수 있습니다. NATS로 다음 작업을 수행할 수 있습니다:
  • 메시지 subject를 게시하거나 구독합니다.
  • 새 메시지가 도착하는 대로 처리합니다.

테이블 생성하기

필수 매개변수:
  • nats_url – host:port(예: localhost:4222)..
  • nats_subjects – NATS 테이블이 구독하거나 게시할 subject 목록입니다. foo.*.bar 또는 baz.> 같은 와일드카드 subject를 지원합니다.
  • nats_format – 메시지 포맷입니다. JSONEachRow와 같이 SQL FORMAT 함수와 동일한 표기법을 사용합니다. 자세한 내용은 포맷 섹션을 참조하십시오.
선택적 매개변수:
  • nats_schema – 포맷에 스키마 정의가 필요한 경우 반드시 사용해야 하는 매개변수입니다. 예를 들어 Cap’n Proto는 스키마 파일의 경로와 루트 schema.capnp:Message 객체 이름이 필요합니다.
  • nats_stream – NATS JetStream에 있는 기존 스트림의 이름입니다.
  • nats_consumer_name – NATS JetStream에 있는 기존 durable pull consumer의 이름입니다.
  • nats_num_consumers – 테이블당 consumer 수입니다. 기본값: 1. NATS core에서만, 단일 consumer의 처리량이 부족한 경우 consumer를 더 지정하십시오.
  • nats_queue_group – NATS subscriber의 큐 그룹 이름입니다. 기본값은 테이블 이름입니다.
  • 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에서 데이터를 플러시하기 위해 폴링으로 수집하는 행 수입니다. 기본값: max_insert_block_size.
  • nats_flush_interval_ms - NATS에서 읽은 데이터를 플러시하는 타임아웃입니다. 기본값: stream_flush_interval_ms.
  • nats_wait_for_flush_interval - true이면 consumer 큐가 비워지는 즉시 종료되는 대신 백그라운드 스트리밍 사이클이 전체 플러시 인터벌(nats_flush_interval_ms, 또는 그 외에는 stream_flush_interval_ms) 동안 열려 있어, 최대 한 플러시 인터벌의 추가 수집 지연 시간을 감수하고 더 많은 메시지를 하나의 블록에 누적할 수 있습니다. 기본값: false(낮은 지연 시간의 drain-and-go 동작).
  • nats_username - NATS 사용자 이름입니다. 서버 구성 파일에 정의된 명명된 컬렉션에 저장된 경우, 쿼리에서 컬렉션의 nats_url 또는 nats_server_list를 재정의할 수 없습니다.
  • nats_password - NATS 비밀번호입니다. 서버 구성 파일에 정의된 명명된 컬렉션에 저장된 경우, 쿼리에서 컬렉션의 nats_url 또는 nats_server_list를 재정의할 수 없습니다.
  • nats_token - NATS 인증 토큰입니다. 서버 구성 파일에 정의된 명명된 컬렉션에 저장된 경우, 쿼리에서 컬렉션의 nats_url 또는 nats_server_list를 재정의할 수 없습니다.
  • nats_credential_file - NATS 자격 증명 파일의 경로입니다. 서버가 자체 권한으로 경로를 열기 때문에, 쿼리에서 nats_url과 nats_server_list를 재정의하지 않는 서버 구성 파일에 정의된 명명된 컬렉션에서만 허용됩니다. 대신 쿼리에서는 파일 내용을 nats_credentials에 전달하십시오.
  • nats_credentials - NATS 자격 증명 내용(사용자 JWT 및 seed가 포함된 .creds 파일과 동일한 페이로드)입니다. 쿼리에서 사용할 수 있는 유일한 표기이므로, <nats_credential_file overridable="false">로 운영자가 해당 경로를 잠그지 않은 한 명명된 컬렉션에서 상속된 nats_credential_file을 대체하며 충돌하지 않습니다. 명명된 컬렉션이 보유한 자격 증명을 제거하기 위해 빈 문자열을 할당할 수는 없습니다.
  • nats_ca_file - NATS 서버 인증서를 검증하는 데 사용하는 신뢰할 수 있는 CA 인증서가 포함된 파일의 경로입니다. nats_secure가 필요합니다. nats_credential_file과 마찬가지로 서버가 자체 권한으로 경로를 열기 때문에, 쿼리에서 nats_url과 nats_server_list를 재정의하지 않는 서버 구성 파일에 정의된 명명된 컬렉션에서만 허용됩니다.
  • nats_client_cert_file - NATS 서버에 제시하는 클라이언트 인증서의 경로입니다. nats_secure 및 nats_client_key_file이 필요합니다. nats_ca_file과 동일한 소스에서 허용됩니다.
  • nats_client_key_file - nats_client_cert_file의 private key 경로입니다. nats_ca_file과 동일한 소스에서 허용됩니다.
  • nats_startup_connect_tries - 시작 시 연결 시도 횟수입니다. 기본값: 5.
  • nats_max_rows_per_message — 행 기반 포맷에서 하나의 NATS 메시지에 기록할 수 있는 최대 행 수입니다. (기본값: 1)
  • nats_commit_on_select - 쿼리 실행 시 메시지를 커밋합니다. JetStream에만 적용되며 core NATS에는 확인 응답이 없습니다. 기본값: 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_file 및 nats_client_key_file로 제공합니다. 세 항목 모두 연산자 설정입니다. 서버 구성 파일에 정의된 명명된 컬렉션에서 가져옵니다. 테이블이 연결될 때 모든 파일을 읽으므로 읽을 수 없거나 형식이 잘못된 파일이 있으면 핸드셰이크 대신 쿼리가 실패합니다. NATS 테이블에 쓰기: 테이블이 단 하나의 subject에서만 읽는 경우, 모든 삽입은 동일한 subject로 게시됩니다. 그러나 테이블이 여러 subject에서 읽는 경우에는 어느 subject로 게시할지 지정해야 합니다. 따라서 여러 subject를 가진 테이블에 삽입할 때는 stream_like_engine_insert_queue를 설정해야 합니다. 테이블이 읽는 subject 중 하나를 선택하여 해당 subject로 데이터를 게시할 수 있습니다. 예시:
nats 관련 설정과 함께 포맷 설정도 추가할 수 있습니다. 예시:
ClickHouse 구성 파일을 사용하여 NATS 서버 구성을 추가할 수 있습니다. 구체적으로 NATS 엔진의 비밀번호를 추가할 수 있습니다:

설명

각 메시지는 한 번만 읽을 수 있으므로(디버깅은 예외) 메시지를 읽는 용도로는 SELECT가 그다지 유용하지 않습니다. 더 실용적인 방법은 materialized view를 사용해 실시간 스레드를 만드는 것입니다. 이를 위해 다음을 수행합니다.
  1. 엔진을 사용해 NATS consumer를 생성하고 이를 데이터 스트림으로 간주합니다.
  2. 원하는 구조로 테이블을 생성합니다.
  3. 엔진의 데이터를 변환해 앞서 생성한 테이블에 저장하는 materialized view를 생성합니다.
MATERIALIZED VIEW가 엔진에 연결되면 백그라운드에서 데이터 수집을 시작합니다. 이렇게 하면 NATS에서 메시지를 지속적으로 받아 SELECT를 사용해 필요한 포맷으로 변환할 수 있습니다. 하나의 NATS 테이블에는 원하는 수만큼 materialized view를 만들 수 있습니다. 이들은 테이블에서 직접 데이터를 읽는 대신 새 레코드(블록 단위)를 받으므로, 여러 테이블에 서로 다른 수준의 상세도(그룹화 및 집계 적용 여부에 따라)로 쓸 수 있습니다. 예시:
스트림 데이터 수신을 중지하거나 변환 로직을 변경하려면 materialized view를 detach하십시오:
ALTER를 사용해 대상 테이블을 변경하려면, 대상 테이블과 뷰의 데이터 사이에 불일치가 생기지 않도록 materialized view를 비활성화하는 것이 좋습니다.

가상 컬럼

  • _subject - NATS 메시지 subject입니다. 데이터 타입: 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 JetStream과 함께 NATS 엔진을 사용하려면 먼저 NATS 스트림과 durable pull consumer를 생성해야 합니다. 이를 위해 NATS CLI 패키지의 nats 유틸리티를 사용할 수 있습니다. 예시는 다음과 같습니다.
스트림과 durable pull consumer를 생성한 후에는 NATS 엔진 테이블을 만들 수 있습니다. 이를 위해 nats&#95;stream, nats&#95;consumer&#95;name, nats&#95;subjects를 초기화해야 합니다:
JetStream 테이블은 최소 1회 전달을 보장합니다. 메시지는 종속 materialized view에 삽입된 후에만 확인 응답되므로, 삽입이 실패하거나 중단된 메시지는 확인 응답되지 않은 상태로 남아 다시 전달됩니다. JetStream을 사용하지 않는 Core NATS는 확인 응답 또는 재생 기능이 없으므로 최대 1회만 전달되며, 중단된 메시지는 손실됩니다.

데이터 내구성

이 섹션은 JetStream에만 적용됩니다. Core NATS는 위에서 설명한 대로 확인 응답이 없고 최대 1회 전달 방식이므로, 확인 응답된 메시지가 유실될 수 있는 윈도우 자체가 없습니다. JetStream 테이블은 삽입된 데이터가 디스크에 기록되기 전에 OS 페이지 캐시가 폐기되면, 이미 소비된 행을 아무런 경고 없이 유실할 수 있습니다. 배치가 종속 materialized view로 전달되면 consumer는 해당 메시지에 확인 응답을 보내고, 이에 따라 스트림은 그 지점을 지나 진행됩니다. 그러나 삽입된 행은 대상 파트가 fsync된 후에야 내구성이 확보되며, 기본적으로 이 작업은 동기적으로 수행되지 않습니다(fsync_after_insert = 0). 확인 응답 이후 대상 파트가 fsync되기 전에 페이지 캐시가 유실되면 메시지가 더 이상 재전달되지 않으므로, 오류 없이 행이 유실되고 count() 값만 작아집니다. 단순히 프로세스를 kill하는 경우에는 이 문제가 드러나지 않습니다. 커널이 페이지 캐시를 유지하다가 결국 디스크에 기록하기 때문입니다. 반면 페이지 캐시가 유실되면 이 문제가 드러나며, 장치 수준의 전원 손실이나 비정상적인 호스트 또는 커널 재설정이 그 예입니다. 권장되는 materialized view 활용 경로(전체 삽입 pipeline이 완료된 후에만 확인 응답이 전송됨)에서는 대상 MergeTree 테이블에 fsync_after_insert = 1(및 fsync_part_directory = 1)을 설정하면 확인 응답이 전송되기 전에 삽입된 파트의 내구성이 확보되므로 이 윈도우가 크게 좁아집니다. 이 설정은 캐스케이딩된 materialized view 대상을 포함하여 배치가 삽입되는 모든 MergeTree 테이블에서 활성화해야 하며, 기본값으로 남아 있는 테이블은 여전히 파트를 유실할 수 있습니다. 비동기 중간 계층은 이 설정만으로는 내구성을 확보하지 못합니다. 예를 들어 Distributed 대상은 ClickHouse Cloud 외부의 기본값인 distributed_foreground_insert = 0일 때 백그라운드로 삽입하므로, 자체적인 내구성 설정이나 동기 삽입이 필요합니다. 또한 이 완화 방안은 nats_commit_on_select = 1을 사용하는 직접적인 INSERT ... SELECT ... FROM <nats_table>에는 적용되지 않습니다. 이 경우 메시지는 대상이 내구성 있는 파트를 기록한 후가 아니라, 읽기가 끝에 도달한 시점에 확인 응답됩니다.
마지막 수정일 2026년 9월 26일