RabbitMQ 테이블 엔진
RabbitMQ 테이블 엔진
이 엔진은 ClickHouse를 RabbitMQ와 통합할 수 있게 해줘요. RabbitMQ를 사용하면 다음을 할 수 있어요.
- 데이터 흐름을 게시하거나 구독하기
- 스트림이 사용 가능해지면 처리하기
출처: 문서
본문
테이블 생성하기
CREATE TABLE [IF NOT EXISTS] [db.]table_name [ON CLUSTER cluster]
(
name1 [type1],
name2 [type2],
...
) ENGINE = RabbitMQ SETTINGS
rabbitmq_host_port = 'host:port' [or rabbitmq_address = 'amqp(s)://guest:guest@localhost/vhost'],
rabbitmq_exchange_name = 'exchange_name',
rabbitmq_format = 'data_format'[,]
[rabbitmq_exchange_type = 'exchange_type',]
[rabbitmq_routing_key_list = 'key1,key2,...',]
[rabbitmq_secure = 0,]
[rabbitmq_schema = '',]
[rabbitmq_num_consumers = N,]
[rabbitmq_num_queues = N,]
[rabbitmq_queue_base = 'queue',]
[rabbitmq_persistent = 0,]
[rabbitmq_skip_broken_messages = N,]
[rabbitmq_max_block_size = N,]
[rabbitmq_flush_interval_ms = N,]
[rabbitmq_queue_settings_list = 'x-dead-letter-exchange=my-dlx,x-max-length=10,x-overflow=reject-publish',]
[rabbitmq_queue_consume = false,]
[rabbitmq_address = '',]
[rabbitmq_vhost = '/',]
[rabbitmq_username = '',]
[rabbitmq_password = '',]
[rabbitmq_commit_on_select = false,]
[rabbitmq_max_rows_per_message = 1,]
[rabbitmq_handle_error_mode = 'default']
필수 매개변수:
rabbitmq_host_port– host:port (예:localhost:5672)rabbitmq_exchange_name– RabbitMQ exchange 이름rabbitmq_format– 메시지 형식. SQLFORMAT함수와 같은 표기법을 사용해요 (예:JSONEachRow). 자세한 내용은 Formats 섹션을 참고해요
선택 매개변수:
rabbitmq_exchange_type– RabbitMQ exchange 유형:direct,fanout,topic,headers,consistent_hash. 기본값:fanoutrabbitmq_routing_key_list– 라우팅 키의 쉼표로 구분된 목록rabbitmq_schema– 형식이 스키마 정의를 요구할 때 반드시 사용해야 하는 매개변수. 예를 들어 Cap'n Proto는 스키마 파일의 경로와 루트schema.capnp:Message객체의 이름을 요구해요rabbitmq_num_consumers– 테이블당 소비자 수. 한 소비자의 처리량이 부족하면 더 많은 소비자를 지정해요. 기본값:1rabbitmq_num_queues– 총 큐 수. 이 숫자를 늘리면 성능이 크게 향상될 수 있어요. 기본값:1rabbitmq_queue_base– 큐 이름에 대한 힌트를 지정해요. 이 설정의 사용 사례는 아래에 설명돼 있어요rabbitmq_persistent– 1(true)로 설정하면 삽입 쿼리에서 delivery mode가 2로 설정돼요 (메시지를 'persistent'로 표시). 기본값:0rabbitmq_skip_broken_messages– 블록당 스키마와 맞지 않는 메시지에 대한 RabbitMQ 메시지 파서 허용치.rabbitmq_skip_broken_messages = N이면 엔진은 파싱할 수 없는 RabbitMQ 메시지 N개를 건너뛰어요 (메시지 하나는 데이터 행 하나와 같음). 기본값:0rabbitmq_max_block_size– RabbitMQ에서 데이터를 flush하기 전에 모은 행 수. 기본값: max_insert_block_sizerabbitmq_flush_interval_ms– RabbitMQ에서 데이터를 flush하는 타임아웃. 기본값: stream_flush_interval_msrabbitmq_queue_settings_list– 큐를 만들 때 RabbitMQ 설정을 지정할 수 있게 해줘요. 사용 가능한 설정:x-max-length,x-max-length-bytes,x-message-ttl,x-expires,x-priority,x-max-priority,x-overflow,x-dead-letter-exchange,x-queue-type.durable설정은 큐에 자동으로 활성화돼요rabbitmq_address– 연결 주소:amqp(s)://user:password@host:port/vhost. 이 설정과rabbitmq_host_port중 하나를 사용해요; 둘 다 설정되면rabbitmq_address가 사용돼요. 그것의 host와 port는 remote_url_allow_hosts에 대해 검사돼요rabbitmq_vhost– RabbitMQ vhost. 기본값:'/'rabbitmq_queue_consume– 사용자 정의 큐를 사용하고 RabbitMQ 설정(exchanges, queues, bindings 선언)을 하지 않아요. 기본값:falserabbitmq_username– RabbitMQ 사용자 이름rabbitmq_password– RabbitMQ 비밀번호reject_unhandled_messages– 오류 발생 시 메시지를 거부해요 (RabbitMQ 부정 확인을 보냄). 이 설정은rabbitmq_queue_settings_list에x-dead-letter-exchange가 정의되어 있으면 자동으로 활성화돼요rabbitmq_commit_on_select– select 쿼리가 만들어질 때 메시지를 커밋해요. 기본값:falserabbitmq_max_rows_per_message— 행 기반 형식에서 하나의 RabbitMQ 메시지에 쓰는 최대 행 수. 기본값:1rabbitmq_empty_queue_backoff_start_ms— RabbitMQ 큐가 비어 있으면 읽기를 재스케줄하기 위한 시작 backoff 지점rabbitmq_empty_queue_backoff_end_ms— RabbitMQ 큐가 비어 있으면 읽기를 재스케줄하기 위한 끝 backoff 지점rabbitmq_empty_queue_backoff_step_ms— RabbitMQ 큐가 비어 있으면 읽기를 재스케줄하기 위한 backoff 단계rabbitmq_handle_error_mode— RabbitMQ 엔진의 오류 처리 방식. 가능한 값: default(메시지 파싱 실패 시 예외 발생), stream(예외 메시지와 원시 메시지가 가상 컬럼_error와_raw_message에 저장), dead_letter_queue(오류 관련 데이터가 system.dead_letter_queue에 저장)
SSL 연결
rabbitmq_host_port 형식에서는 TLS를 사용하려면 rabbitmq_secure = 1을 설정해요. rabbitmq_address 형식에서는 전송(transport)이 URI 스킴에서 나오므로 amqps를 사용해요: rabbitmq_address = 'amqps://guest:guest@localhost/vhost'. rabbitmq_secure는 주소 형식에서 무시되고, 일반 텍스트 amqp:// 주소와 함께 rabbitmq_secure = 1은 조용히 평문으로 연결하는 대신 거부돼요.
사용된 라이브러리의 기본 동작은 생성된 TLS 연결이 충분히 안전한지 확인하지 않는 것이에요. 인증서가 만료, 자체 서명, 누락 또는 무효인지와 상관없이: 연결은 단순히 허용돼요. 더 엄격한 인증서 검사는 향후 구현될 수 있어요.
또한 rabbitmq 관련 설정과 함께 형식 설정을 추가할 수 있어요.
예시:
CREATE TABLE queue (
key UInt64,
value UInt64,
date DateTime
) ENGINE = RabbitMQ SETTINGS rabbitmq_host_port = 'localhost:5672',
rabbitmq_exchange_name = 'exchange1',
rabbitmq_format = 'JSONEachRow',
rabbitmq_num_consumers = 5,
date_time_input_format = 'best_effort';
RabbitMQ 서버 구성은 ClickHouse 구성 파일로 추가해야 해요.
필수 구성:
<rabbitmq>
<username>root</username>
<password>clickhouse</password>
</rabbitmq>
추가 구성:
<rabbitmq>
<vhost>clickhouse</vhost>
</rabbitmq>
설명
SELECT는 (디버깅 외에는) 메시지 읽기에 특히 유용하지 않아요. 각 메시지는 한 번만 읽을 수 있기 때문이에요. materialized view를 사용해 실시간 스레드를 만드는 것이 더 실용적이에요. 그러려면:
- 엔진을 사용해 RabbitMQ 소비자를 만들고 그것을 데이터 스트림으로 간주해요
- 원하는 구조로 테이블을 만들어요
- 엔진의 데이터를 변환해 앞서 만든 테이블에 넣는 materialized view를 만들어요
MATERIALIZED VIEW가 엔진에 조인되면 백그라운드에서 데이터 수집을 시작해요. 이렇게 하면 RabbitMQ에서 메시지를 계속 받아 SELECT로 필요한 형식으로 변환할 수 있어요. 하나의 RabbitMQ 테이블은 원하는 만큼 많은 materialized view를 가질 수 있어요.
데이터는 rabbitmq_exchange_type과 지정된 rabbitmq_routing_key_list에 따라 채널로 나눌 수 있어요. 테이블당 하나 이상의 exchange를 가질 수 없어요. 하나의 exchange는 여러 테이블 간에 공유될 수 있어요 — 동시에 여러 테이블로의 라우팅을 가능하게 해요.
Exchange 유형 옵션:
direct— 키의 정확한 일치를 기반으로 라우팅돼요. 예: 테이블 키 목록key1,key2,key3,key4,key5, 메시지 키가 그중 하나와 같을 수 있어요fanout— 키와 관계없이 (exchange 이름이 같은) 모든 테이블로 라우팅돼요topic— 점으로 구분된 키가 있는 패턴을 기반으로 라우팅돼요. 예:*.logs,records.*.*.2020,*.2018,*.2019,*.2020headers—x-match=all또는x-match=any설정과 함께key=value일치를 기반으로 라우팅돼요. 예: 테이블 키 목록x-match=all,format=logs,type=report,year=2020consistent_hash— (exchange 이름이 같은) 모든 바인딩 테이블 간에 데이터가 고르게 분배돼요. 이 exchange 유형은 RabbitMQ 플러그인으로 활성화해야 한다는 점을 유의해요:rabbitmq-plugins enable rabbitmq_consistent_hash_exchange
rabbitmq_queue_base 설정은 다음 경우에 사용될 수 있어요.
- 서로 다른 테이블이 큐를 공유할 수 있게 해서, 여러 소비자가 같은 큐에 등록될 수 있어 더 나은 성능을 만들어요.
rabbitmq_num_consumers및/또는rabbitmq_num_queues설정을 사용하면 이 매개변수들이 같으면 큐의 정확한 일치가 달성돼요 - 모든 메시지가 성공적으로 소비되지 않았을 때 특정 durable 큐에서 읽기를 복원할 수 있게 해요. 특정 큐 하나에서 소비를 재개하려면
rabbitmq_queue_base설정에 그 이름을 설정하고rabbitmq_num_consumers와rabbitmq_num_queues를 지정하지 않아요 (기본값 1). 특정 테이블에 대해 선언된 모든 큐에서 소비를 재개하려면 같은 설정들(rabbitmq_queue_base,rabbitmq_num_consumers,rabbitmq_num_queues)을 지정하기만 하면 돼요. 기본적으로 큐 이름은 테이블마다 고유해요 - 큐가 durable로 선언되고 자동 삭제되지 않으므로 재사용하려면. (RabbitMQ CLI 도구 중 아무거나로 삭제할 수 있어요)
성능을 향상시키기 위해 수신된 메시지는 max_insert_block_size 크기의 블록으로 그룹화돼요. 블록이 stream_flush_interval_ms 밀리초 안에 형성되지 않으면 블록의 완전성과 관계없이 데이터가 테이블로 flush돼요.
rabbitmq_exchange_type과 함께 rabbitmq_num_consumers 및/또는 rabbitmq_num_queues 설정이 지정되면:
rabbitmq-consistent-hash-exchange플러그인이 활성화되어야 해요- 게시된 메시지의
message_id속성이 지정되어야 해요 (각 메시지/배치에 대해 고유)
삽입 쿼리의 경우 각 게시된 메시지에 추가되는 메시지 메타데이터가 있어요: messageID와 republished 플래그(두 번 이상 게시되면 true) — 메시지 헤더로 접근할 수 있어요.
삽입과 materialized views에 같은 테이블을 사용하지 마세요.
예시:
CREATE TABLE queue (
key UInt64,
value UInt64
) ENGINE = RabbitMQ SETTINGS rabbitmq_host_port = 'localhost:5672',
rabbitmq_exchange_name = 'exchange1',
rabbitmq_exchange_type = 'headers',
rabbitmq_routing_key_list = 'format=logs,type=report,year=2020',
rabbitmq_format = 'JSONEachRow',
rabbitmq_num_consumers = 5;
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;
가상 컬럼
_exchange_name- RabbitMQ exchange 이름. 데이터 타입:String_channel_id- 메시지를 받은 소비자가 선언된 ChannelID. 데이터 타입:String_delivery_tag- 수신된 메시지의 DeliveryTag. 채널마다 범위가 정해짐. 데이터 타입:UInt64_redelivered- 메시지의redelivered플래그. 데이터 타입:UInt8_message_id- 수신된 메시지의 messageID; 메시지가 게시될 때 설정되었으면 비어 있지 않음. 데이터 타입:String_timestamp- 수신된 메시지의 타임스탬프; 메시지가 게시될 때 설정되었으면 비어 있지 않음. 데이터 타입:UInt64
rabbitmq_handle_error_mode='stream'일 때 추가 가상 컬럼:
_raw_message- 성공적으로 파싱되지 못한 원시 메시지. 데이터 타입:Nullable(String)_error- 실패한 파싱 중 발생한 예외 메시지. 데이터 타입:Nullable(String)
참고: _raw_message와 _error 가상 컬럼은 파싱 중 예외가 발생한 경우에만 채워지며, 메시지가 성공적으로 파싱되면 항상 NULL이에요.
주의사항(Caveats)
테이블 정의에 기본 컬럼 표현식(DEFAULT, MATERIALIZED, ALIAS 같은)을 지정해도 무시돼요. 대신 컬럼들은 타입에 해당하는 기본값으로 채워져요.
데이터 형식 지원
RabbitMQ 엔진은 ClickHouse에서 지원하는 모든 형식을 지원해요. 하나의 RabbitMQ 메시지의 행 수는 형식이 행 기반인지 블록 기반인지에 따라 달라져요.
- 행 기반 형식의 경우 하나의 RabbitMQ 메시지의 행 수를
rabbitmq_max_rows_per_message설정으로 제어할 수 있어요 - 블록 기반 형식의 경우 블록을 더 작은 부분으로 나눌 수 없지만, 하나의 블록의 행 수는 일반 설정 max_block_size로 제어할 수 있어요
전원 손실 시 데이터 내구성(Durability)
RabbitMQ 엔진은 삽입된 데이터가 디스크에 쓰이기 전에 OS 페이지 캐시가 버려지면 이미 소비된 행을 조용히 잃을 수 있어요. 배치가 의존하는 materialized view에 푸시된 후 소비자가 브로커에 basic.ack을 보내 브로커가 그 메시지들을 삭제할 수 있게 해요. 그러나 삽입된 행은 대상 파트가 fsync된 후에만 내구성이 생기는데, 기본적으로 동기적으로 발생하지 않아요 (fsync_after_insert = 0). 확인(acknowledgement) 후 대상 파트가 fsync되기 전에 페이지 캐시가 손실되면 브로커가 이미 메시지를 버렸고 소비자가 재연결 시 그 메시지들을 넘어 재개하므로 행이 오류 없이 손실되고 count()가 단순히 더 작아져요. 일반적인 프로세스 종료는 커널이 페이지 캐시를 유지하고 결국 기록하므로 이를 드러내지 않아요. 페이지 캐시 손실은 드러내요. 예로 장치 수준 전원 손실과 비정상적인 호스트나 커널 리셋이 있어요.
권장하는 materialized-view 소비 경로(전체 삽입 파이프라인이 끝난 뒤에만 확인을 보내는)에서는 대상 MergeTree 테이블에 fsync_after_insert = 1(그리고 fsync_part_directory = 1)을 설정하면 삽입된 파트가 확인 전에 내구성이 생겨 이 창이 크게 좁아져요. 이 설정은 배치가 삽입되는 모든 MergeTree 테이블, 캐스케이드된 materialized-view 대상까지 포함해 활성화해야 해요. 기본값으로 남겨둔 테이블은 여전히 파트를 잃을 수 있어요. 비동기 중개자는 이 설정만으로 내구성을 얻지 못해요. 예를 들어 distributed_foreground_insert = 0이면 Distributed 대상이 백그라운드에서 삽입하는데, 이는 ClickHouse Cloud 외부의 기본값이므로 자체 내구성 설정이나 동기 삽입이 필요해요. 이 완화는 rabbitmq_commit_on_select = 1을 사용한 직접 INSERT ... SELECT ... FROM <rabbitmq_table>에도 적용되지 않아요. 그 경우 메시지는 대상이 내구성 있는 파트를 쓴 뒤가 아니라 읽기가 끝에 도달할 때 확인되기 때문이에요.