Distributed 테이블 엔진
Distributed 테이블 엔진
Distributed 엔진은 자체적으로 데이터를 저장하지 않지만, 여러 서버에서 분산 쿼리 처리를 가능하게 하는 테이블 엔진이에요. 읽기는 자동으로 병렬화되며 원격 서버의 인덱스도 활용돼요. ClickHouse Cloud에서는 Remote 또는 remoteSecure 테이블 함수를 사용해야 해요.
출처: 문서
본문
ClickHouse Cloud에서 분산 테이블 엔진을 만들려면 remote와 remoteSecure 테이블 함수를 사용할 수 있어요. Distributed(...) 구문은 ClickHouse Cloud에서 사용할 수 없습니다.
Distributed 엔진을 가진 테이블은 자체적으로는 데이터를 저장하지 않지만, 여러 서버에서 분산 쿼리 처리를 가능하게 해요. 읽기는 자동으로 병렬화됩니다. 읽는 동안 원격 서버에 테이블 인덱스가 있다면 그것을 사용해요.
테이블 생성 (Creating a table)
CREATE TABLE [IF NOT EXISTS] [db.]table_name [ON CLUSTER cluster]
(
name1 [type1] [DEFAULT|MATERIALIZED|ALIAS expr1],
name2 [type2] [DEFAULT|MATERIALIZED|ALIAS expr2],
...
) ENGINE = Distributed(cluster, database, table[, sharding_key[, policy_name]])
[SETTINGS name=value, ...]
테이블에서 생성 (From a table)
Distributed 테이블이 현재 서버의 테이블을 가리킬 때, 해당 테이블의 스키마를 채택할 수 있어요:
CREATE TABLE [IF NOT EXISTS] [db.]table_name [ON CLUSTER cluster] AS [db2.]name2 ENGINE = Distributed(cluster, database, table[, sharding_key[, policy_name]]) [SETTINGS name=value, ...]
Remote 및 RemoteSecure 엔진 (Remote and RemoteSecure engines)
Remote와 RemoteSecure는 remote 및 remoteSecure 테이블 함수와 동일한 주소 표현식과 자격 증명을 받아들이는 영구 테이블 엔진이에요:
CREATE TABLE [IF NOT EXISTS] [db.]table_name
(
name1 [type1],
name2 [type2],
...
) ENGINE = Remote(addresses_expr, [db, table, [user [, password], sharding_key]])
[SETTINGS name = value, ...]
RemoteSecure는 동일한 인자를 받아들이며 보안 연결로 접속합니다(기본적으로 보안 TCP 포트를 사용해요). 인자는 remote 및 remoteSecure 테이블 함수와 정확히 동일하게 해석되며, 지원되는 시그니처에 대한 설명은 해당 함수 설명을 참조하세요. 테이블 구조는 생략할 수 있으며, 이 경우 원격 테이블에서 추론됩니다.
생성된 스토리지의 설정(예: skip_unavailable_shards)은 엔진 정의 뒤에 지정돼요. 예를 들어 ENGINE = Remote('127.0.0.1', system, one) SETTINGS skip_unavailable_shards = 1처럼요. 참고로 remote 및 remoteSecure 테이블 함수는 인자 중에 SETTINGS 절을 받아들이는데(remote('127.0.0.1', system.one, SETTINGS skip_unavailable_shards = 1)), 테이블 함수에게는 둘 곳이 없기 때문이에요. 즉 엔진은 그 형태를 받아들이지 않습니다.
예를 들어:
CREATE TABLE remote_one ENGINE = Remote('127.0.0.1', system, one);
SELECT * FROM remote_one;
이것은 CREATE TABLE ... AS remote(...)의 영구적 등가물입니다. remote 테이블 함수와 마찬가지로 이 엔진들은 편리하지만, 구성된 클러스터 위에 Distributed를 쓰는 것처럼 샤드와 레플리카를 선언적으로 설정할 수는 없어요. 따라서 영구적이고 자주 사용하는 서버 집합이라면 클러스터를 정의하고 Distributed 엔진을 사용하는 것이 좋습니다.
대상은 테이블 함수일 수도 있는데, 예를 들어 Remote('127.0.0.1', numbers(10)) 또는 Remote('127.0.0.1', merge(db, '^table_'))처럼요. 이런 테이블은 읽기 전용입니다. 삽입할 원격 테이블이 없으므로 INSERT는 NOT_IMPLEMENTED 예외로 거부돼요. 이 읽기 전용 제한은 remote 및 remoteSecure 테이블 함수에도 동일하게 적용됩니다. 일반적인 db/table 대상에서는 SELECT와 INSERT가 모두 지원되지만, 테이블 함수 대상(remote('127.0.0.1', numbers(10)))은 같은 이유로 읽기 전용이에요.
Distributed 매개변수 (Distributed parameters)
| Parameter | Description |
|---|---|
| cluster | 서버 설정 파일에 있는 클러스터 이름 |
| database | 원격 데이터베이스의 이름 |
| table | 원격 테이블의 이름 |
| sharding_key (선택) | 샤딩 키. sharding_key 지정은 다음 경우에 필요합니다: 분산 테이블에 대한 INSERT 시 (테이블 엔진이 데이터를 어떻게 분할할지 판단하기 위해 sharding_key가 필요하기 때문). 단, insert_distributed_one_random_shard 설정이 활성화되어 있으면 INSERT에 sharding 키가 필요하지 않아요. optimize_skip_unused_shards 사용 시, 어떤 샤드를 쿼리해야 하는지 결정하는 데 sharding_key가 필요하므로 사용합니다 |
| policy_name (선택) | 배경 전송용 임시 파일을 저장하는 데 사용되는 정책 이름 |
Distributed 설정 (Distributed settings)
| Setting | Description | Default value |
|---|---|---|
| fsync_after_insert | Distributed에 대한 배경 INSERT 후 파일 데이터에 대한 fsync를 수행. OS가 삽입된 전체 데이터를 초기자 노드 디스크의 파일로 플러시했음을 보장. | false |
| fsync_directories | 디렉토리에 대한 fsync 수행. Distributed 테이블의 배경 INSERT와 관련된 연산(삽입 후, 데이터를 샤드로 보낸 후 등) 후 OS가 디렉토리 메타데이터를 갱신했음을 보장. | false |
| skip_unavailable_shards | true이면 ClickHouse가 사용할 수 없는 샤드를 조용히 건너뜁니다. 이 설정의 동작은 skip_unavailable_shards_mode 매개변수에 의해 제어됩니다. | false |
| skip_unavailable_shards_mode | skip_unavailable_shards가 활성화됐을 때 원격 샤드의 어떤 예외를 무시할지 제어합니다: unavailable은 연결 오류만 무시하고, unavailable_or_table_missing은 테이블이나 데이터베이스 누락도 무시하며, unavailable_or_exception_before_processing은 샤드가 데이터를 반환하기 전에 발생한 모든 예외를 무시합니다. | unavailable_or_table_missing |
| bytes_to_throw_insert | 이 숫자보다 많은 압축 바이트가 배경 INSERT를 위해 대기 중이면 예외가 발생합니다. 0 - 발생시키지 않음. | 0 |
| bytes_to_delay_insert | 이 숫자보다 많은 압축 바이트가 배경 INSERT를 위해 대기 중이면 쿼리가 지연됩니다. 0 - 지연시키지 않음. | 0 |
| max_delay_to_insert | 배경 전송을 위해 대기 중인 바이트가 많을 때 Distributed 테이블에 데이터 삽입의 최대 지연(초). | 60 |
| background_insert_batch | distributed_background_insert_batch와 동일 | 0 |
| background_insert_split_batch_on_failure | distributed_background_insert_split_batch_on_failure와 동일 | 0 |
| background_insert_sleep_time_ms | distributed_background_insert_sleep_time_ms와 동일 | 0 |
| background_insert_max_sleep_time_ms | distributed_background_insert_max_sleep_time_ms와 동일 | 0 |
| flush_on_detach | DETACH / DROP / 서버 종료 시 원격 노드로 데이터를 플러시. | true |
내구성(durability) 설정(fsync_...) 관련:
- 배경
INSERT(즉distributed_foreground_insert=false)에만 영향을 줍니다. 이때 데이터는 먼저 초기자 노드 디스크에 저장된 뒤, 나중에 배경에서 샤드로 전송됩니다. INSERT성능을 크게 저하시킬 수 있어요- Distributed 테이블 폴더 안에 저장된 데이터를 당신의 INSERT를 받아들인 노드에 기록하는 데 영향을 줍니다. 기본 MergeTree 테이블에 데이터 기록을 보장해야 한다면
system.merge_tree_settings의 내구성 설정(...fsync...)을 참고하세요
Insert 제한 설정(..._insert) 관련해서는 다음도 참고하세요:
- distributed_foreground_insert 설정
- prefer_localhost_replica 설정
bytes_to_throw_insert는bytes_to_delay_insert보다 먼저 처리되므로,bytes_to_delay_insert보다 작은 값으로 설정하면 안 됩니다
예시
CREATE TABLE hits_all AS hits
ENGINE = Distributed(logs, default, hits[, sharding_key[, policy_name]])
SETTINGS
fsync_after_insert=0,
fsync_directories=0;
logs 클러스터의 모든 서버, 클러스터에 있는 모든 서버의 default.hits 테이블에서 데이터를 읽습니다. 데이터는 읽기만 되는 것이 아니라 원격 서버에서 부분적으로 처리됩니다(가능한 범위까지). 예를 들어 GROUP BY가 있는 쿼리에서는 원격 서버에서 데이터가 집계되고, 집계 함수의 중간 상태가 요청자 서버로 전송됩니다. 그 다음 데이터가 추가로 집계됩니다.
데이터베이스 이름 대신 문자열을 반환하는 상수 표현식을 사용할 수 있어요. 예: currentDatabase().
클러스터 (Clusters)
클러스터는 서버 설정 파일에 구성됩니다:
<remote_servers>
<logs>
<!-- Inter-server per-cluster secret for Distributed queries
default: no secret (no authentication will be performed)
If set, then Distributed queries will be validated on shards, so at least:
- such cluster should exist on the shard,
- such cluster should have the same secret.
And also (and which is more important), the initial_user will
be used as current user for the query.
-->
<!-- <secret></secret> -->
<!-- Optional. Whether distributed DDL queries (ON CLUSTER clause) are allowed for this cluster. Default: true (allowed). -->
<!-- <allow_distributed_ddl_queries>true</allow_distributed_ddl_queries> -->
<shard>
<!-- Optional. Shard weight when writing data. Default: 1. -->
<weight>1</weight>
<!-- Optional. The shard name. Must be non-empty and unique among shards in the cluster. If not specified, will be empty. -->
<name>shard_01</name>
<!-- Optional. Whether to write data to just one of the replicas. Default: false (write data to all replicas). -->
<internal_replication>false</internal_replication>
<replica>
<!-- Optional. Priority of the replica for load balancing (see also load_balancing setting). Default: 1 (less value has more priority). -->
<priority>1</priority>
<host>example01-01-1</host>
<port>9000</port>
</replica>
<replica>
<host>example01-01-2</host>
<port>9000</port>
</replica>
</shard>
<shard>
<weight>2</weight>
<name>shard_02</name>
<internal_replication>false</internal_replication>
<replica>
<host>example01-02-1</host>
<port>9000</port>
</replica>
<replica>
<host>example01-02-2</host>
<secure>1</secure>
<port>9440</port>
</replica>
</shard>
</logs>
</remote_servers>
여기서 logs라는 이름의 클러스터가 정의되었고, 두 개의 샤드로 구성되며 각각 두 개의 레플리카를 포함합니다. 샤드는 데이터의 서로 다른 부분을 담고 있는 서버를 가리킵니다(모든 데이터를 읽으려면 모든 샤드에 접근해야 해요). 레플리카는 중복되는 서버입니다(모든 데이터를 읽으려면 레플리카 중 하나에만 접근하면 됩니다).
클러스터 이름에는 점(dot)이 포함되면 안 됩니다.
각 서버에 대해 host, port, 그리고 선택적으로 user, password, secure, compression, bind_host 매개변수가 지정됩니다:
| Parameter | Description | Default Value |
|---|---|---|
| host | 원격 서버의 주소. 도메인 또는 IPv4, IPv6 주소를 사용할 수 있어요. 도메인을 지정하면 서버가 시작할 때 DNS 요청을 하고, 그 결과는 서버가 실행되는 동안 저장됩니다. DNS 요청이 실패하면 서버가 시작되지 않아요. DNS 레코드를 변경하면 서버를 다시 시작하세요. | - |
| port | 메신저 활동을 위한 TCP 포트(설정에서 tcp_port, 보통 9000으로 설정). http_port와 혼동하지 마세요. | - |
| user | 원격 서버에 연결하기 위한 사용자 이름. 이 사용자는 지정된 서버에 연결할 수 있는 액세스 권한이 있어야 해요. 액세스는 users.xml 파일에서 구성됩니다. 자세한 내용은 Access rights 섹션을 참조하세요. | default |
| password | 원격 서버에 연결하기 위한 비밀번호(마스킹되지 않음). | " |
| secure | 보안 SSL/TLS 연결을 사용할지 여부. 보통 포트도 지정해야 해요(기본 보안 포트는 9440). 서버는 <tcp_port_secure>9440</tcp_port_secure>에서 수신하고 올바른 인증서로 구성되어야 합니다. | false |
| compression | 데이터 압축 사용. | true |
| bind_host | 이 노드에서 원격 서버에 연결할 때 사용할 소스 주소. IPv4 주소만 지원됩니다. ClickHouse 분산 쿼리가 사용하는 소스 IP 주소 설정이 필요한 고급 배포 사례를 위한 것입니다. | - |
레플리카를 지정하면 읽을 때 각 샤드에 대해 사용 가능한 레플리카 중 하나가 선택됩니다. 로드 밸런싱 알고리즘(어느 레플리카에 접근할지에 대한 선호도)을 구성할 수 있어요 – load_balancing 설정을 참조하세요. 서버와의 연결이 수립되지 않으면 짧은 타임아웃으로 연결을 시도합니다. 연결이 실패하면 다음 레플리카가 선택되고, 모든 레플리카에 대해 이런 식으로 계속됩니다. 모든 레플리카에 대한 연결 시도가 실패하면 같은 방식으로 시도를 여러 번 반복합니다. 이것은 복원력에 유리하지만 완전한 내결함성을 제공하지는 않아요. 원격 서버가 연결을 수락할 수도 있지만 작동하지 않거나 제대로 작동하지 않을 수 있기 때문입니다.
샤드 하나만 지정할 수도 있으며(이 경우 쿼리 처리는 distributed가 아니라 remote라고 불러야 해요), 원하는 수의 샤드까지 지정할 수 있습니다. 각 샤드에는 1개부터 원하는 수의 레플리카를 지정할 수 있어요. 각 샤드마다 서로 다른 수의 레플리카를 지정할 수 있습니다.
설정에 원하는 만큼 많은 클러스터를 지정할 수 있어요. 클러스터를 보려면 system.clusters 테이블을 사용하세요.
Distributed 엔진을 사용하면 클러스터를 로컬 서버처럼 다룰 수 있어요. 그러나 클러스터의 구성은 동적으로 지정할 수 없고, 서버 설정 파일에 구성해야 합니다. 일반적으로 클러스터의 모든 서버는 동일한 클러스터 구성을 갖습니다(필수는 아니지만요). 설정 파일의 클러스터는 서버를 재시작하지 않고도 즉시 업데이트됩니다.
매번 알 수 없는 샤드와 레플리카 집합에 쿼리를 보내야 한다면 Distributed 테이블을 만들 필요가 없어요. 대신 remote 테이블 함수를 사용하세요. Table functions 섹션을 참조하세요.
데이터 쓰기 (Writing data)
클러스터에 데이터를 쓰는 방법은 두 가지가 있어요.
첫째, 어떤 서버에 어떤 데이터를 쓸지 정의하고 각 샤드에 직접 쓰기를 수행할 수 있어요. 즉 Distributed 테이블이 가리키는 클러스터의 원격 테이블에 직접 INSERT 문을 수행하는 것입니다. 영역의 요구사항 때문에 사소하지 않은 샤딩 스킴도 사용할 수 있으므로 가장 유연한 해결책입니다. 데이터를 각기 다른 샤드에 완전히 독립적으로 쓸 수 있으므로 가장 최적의 해결책이기도 해요.
둘째, Distributed 테이블에 INSERT 문을 수행할 수 있어요. 이 경우 테이블이 삽입된 데이터를 서버 간에 분배합니다. Distributed 테이블에 쓰려면 sharding_key 매개변수가 구성되어 있어야 합니다(샤드가 하나뿐인 경우는 제외).
같은 클러스터를 사용하는 Distributed 테이블 간의 호환 가능한 INSERT ... SELECT 쿼리의 경우, parallel_distributed_insert_select가 각 샤드에서 쿼리를 병렬로 실행할 수 있어요.
각 샤드는 설정 파일에 <weight>를 가질 수 있어요. 기본적으로 weight는 1입니다. 데이터는 샤드 weight에 비례하는 양으로 샤드 간에 분배됩니다. 모든 샤드 weight가 합산된 다음, 각 샤드의 weight를 총계로 나누어 각 샤드의 비율을 결정합니다. 예를 들어 샤드가 두 개이고 첫 번째 weight가 1, 두 번째 weight가 2라면 첫 번째는 삽입된 행의 3분의 1(1 / 3), 두 번째는 3분의 2(2 / 3)를 받습니다.
각 샤드는 설정 파일에 internal_replication 매개변수를 가질 수 있어요. 이 매개변수가 true로 설정되면 쓰기 연산이 첫 번째 정상 레플리카를 선택하여 데이터를 씁니다. Distributed 테이블의 기반이 되는 테이블이 레플리케이션 테이블(예: Replicated*MergeTree 테이블 엔진)인 경우 이를 사용하세요. 테이블 레플리카 중 하나가 쓰기를 받고, 다른 레플리카로 자동으로 레플리케이션됩니다.
internal_replication이 false로 설정되면(기본값) 데이터가 모든 레플리카에 기록됩니다. 이 경우 Distributed 테이블이 자체적으로 데이터를 레플리케이션합니다. 레플리카의 일관성이 검사되지 않고 시간이 지나면 약간 다른 데이터를 포함하게 되므로, 레플리케이션된 테이블을 사용하는 것보다 좋지 않아요.
데이터 행이 전송될 샤드를 선택하기 위해 샤딩 표현식을 분석하고, 이를 샤드의 총 weight로 나눈 나머지를 취합니다. 행은 prev_weights에서 prev_weights + weight 사이의 반열린 구간에 해당하는 샤드로 전송됩니다. 여기서 prev_weights는 가장 작은 번호를 가진 샤드들의 총 weight이고, weight는 이 샤드의 weight입니다. 예를 들어 샤드가 두 개이고 첫 번째 weight가 9, 두 번째 weight가 10이라면 나머지 [0, 9) 범위에 대해 첫 번째 샤드로, 나머지 [9, 19) 범위에 대해 두 번째 샤드로 행이 전송됩니다.
샤딩 표현식은 정수를 반환하는 상수와 테이블 컬럼의 어떤 표현식이든 될 수 있습니다. 예를 들어 데이터의 무작위 분배를 위해 rand() 표현식을, 사용자 ID를 나눈 나머지로 분배하기 위해 UserID를 사용할 수 있어요(그러면 단일 사용자의 데이터가 단일 샤드에 위치하게 되어 사용자 기준 IN과 JOIN 실행이 단순해집니다). 컬럼 중 하나가 충분히 균등하게 분배되지 않으면 intHash64(UserID) 같은 해시 함수로 감쌀 수 있어요.
단순한 나눗셈 나머지는 샤딩의 제한된 해결책이며 항상 적절하지는 않아요. 중간 및 대용량 데이터(수십 대 서버)에서는 잘 작동하지만, 초대용량 데이터(수백 대 이상의 서버)에서는 적절하지 않습니다. 후자의 경우에는 Distributed 테이블 항목을 사용하지 말고 영역 요구사항에 따른 샤딩 스킴을 사용하세요.
다음과 같은 경우 샤딩 스킴을 고려해야 해요:
- 특정 키로 데이터를 조인해야 하는 쿼리(
IN또는JOIN)를 사용하는 경우. 데이터가 이 키로 샤딩되면 GLOBAL IN 또는 GLOBAL JOIN 대신 로컬 IN 또는 JOIN을 사용할 수 있는데, 훨씬 더 효율적입니다. - 대량의 서버(수백 대 이상)를 소수의 작은 쿼리와 함께 사용하는 경우. 예를 들어 개별 클라이언트(웹사이트, 광고주, 파트너 등)의 데이터에 대한 쿼리처럼요. 작은 쿼리가 전체 클러스터에 영향을 주지 않도록 하려면 단일 클라이언트의 데이터를 단일 샤드에 두는 것이 합리적입니다. 또는 2단계 샤딩을 설정할 수 있어요. 전체 클러스터를 "레이어"로 나누고, 레이어는 여러 샤드로 구성됩니다. 단일 클라이언트의 데이터는 단일 레이어에 위치하지만, 필요에 따라 레이어에 샤드를 추가할 수 있고 그 안에서는 데이터가 무작위로 분배됩니다. 각 레이어에 대해
Distributed테이블이 생성되고, 전역 쿼리를 위해 단일 공유 분산 테이블이 생성됩니다.
데이터는 배경에서 기록됩니다. 테이블에 삽입하면 데이터 블록이 로컬 파일 시스템에만 기록됩니다. 데이터는 가능한 한 빨리 배경에서 원격 서버로 전송됩니다. 데이터 전송 주기는 distributed_background_insert_sleep_time_ms와 distributed_background_insert_max_sleep_time_ms 설정으로 관리됩니다. Distributed 엔진은 삽입된 데이터가 있는 각 파일을 개별적으로 전송하지만, distributed_background_insert_batch 설정으로 파일의 배치 전송을 활성화할 수 있어요. 이 설정은 로컬 서버와 네트워크 리소스를 더 잘 활용하여 클러스터 성능을 향상시킵니다. 테이블 디렉토리의 파일 목록(전송 대기 중인 데이터)을 확인하여 데이터가 성공적으로 전송되었는지 확인해야 합니다: /var/lib/clickhouse/data/database/table/. 배경 작업을 수행하는 스레드 수는 background_distributed_schedule_pool_size 설정으로 설정할 수 있어요.
서버가 사라지거나 Distributed 테이블에 대한 INSERT 후 거친 재시작(예: 하드웨어 오류로 인해)을 겪으면 삽입된 데이터가 손실될 수 있어요. 테이블 디렉토리에서 손상된 데이터 part가 감지되면 broken 하위 디렉토리로 이동되어 더 이상 사용되지 않습니다.
데이터 읽기 (Reading data)
Distributed 테이블을 쿼리하면 SELECT 쿼리가 모든 샤드로 전송되며, 데이터가 샤드 간에 어떻게 분배되든(완전히 무작위로 분배될 수도 있음) 상관없이 작동합니다. 새 샤드를 추가할 때 기존 데이터를 그 샤드로 옮길 필요는 없어요. 대신 더 큰 weight를 사용해 새 데이터를 그 샤드에 쓸 수 있습니다. 데이터가 약간 고르지 않게 분배되겠지만 쿼리는 정확하고 효율적으로 작동합니다.
max_parallel_replicas 옵션이 활성화되면 단일 샤드 내의 모든 레플리카에서 쿼리 처리가 병렬화됩니다. 자세한 내용은 max_parallel_replicas 섹션을 참조하세요.
distributed in 및 global in 쿼리가 어떻게 처리되는지 자세히 알아보려면 이 문서를 참조하세요.
가상 컬럼 (Virtual columns)
_Shard_num
_shard_num — system.clusters 테이블의 shard_num 값을 포함합니다. 타입: UInt32.
remote 및 cluster 테이블 함수는 내부적으로 임시 Distributed 테이블을 만들기 때문에 _shard_num도 거기에서 사용할 수 있어요.