NATS로 다음 작업을 수행할 수 있습니다:
- 메시지 subject를 게시하거나 구독합니다.
- 새 메시지가 도착하는 대로 처리합니다.
테이블 생성하기
nats_url– host:port(예:localhost:4222)..nats_subjects– NATS 테이블이 구독하거나 게시할 subject 목록입니다.foo.*.bar또는baz.>같은 와일드카드 subject를 지원합니다.nats_format– 메시지 포맷입니다.JSONEachRow와 같이 SQLFORMAT함수와 동일한 표기법을 사용합니다. 자세한 내용은 포맷 섹션을 참조하십시오.
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에 저장됨).
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로 데이터를 게시할 수 있습니다. 예시:
설명
각 메시지는 한 번만 읽을 수 있으므로(디버깅은 예외) 메시지를 읽는 용도로는SELECT가 그다지 유용하지 않습니다. 더 실용적인 방법은 materialized view를 사용해 실시간 스레드를 만드는 것입니다. 이를 위해 다음을 수행합니다.
- 엔진을 사용해 NATS consumer를 생성하고 이를 데이터 스트림으로 간주합니다.
- 원하는 구조로 테이블을 생성합니다.
- 엔진의 데이터를 변환해 앞서 생성한 테이블에 저장하는 materialized view를 생성합니다.
MATERIALIZED VIEW가 엔진에 연결되면 백그라운드에서 데이터 수집을 시작합니다. 이렇게 하면 NATS에서 메시지를 지속적으로 받아 SELECT를 사용해 필요한 포맷으로 변환할 수 있습니다.
하나의 NATS 테이블에는 원하는 수만큼 materialized view를 만들 수 있습니다. 이들은 테이블에서 직접 데이터를 읽는 대신 새 레코드(블록 단위)를 받으므로, 여러 테이블에 서로 다른 수준의 상세도(그룹화 및 집계 적용 여부에 따라)로 쓸 수 있습니다.
예시:
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 생성
durable pull consumer 생성
nats_stream, nats_consumer_name, nats_subjects를 초기화해야 합니다:
데이터 내구성
이 섹션은 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>에는 적용되지 않습니다. 이 경우 메시지는 대상이 내구성 있는 파트를 기록한 후가 아니라, 읽기가 끝에 도달한 시점에 확인 응답됩니다.