> ## Documentation Index
> Fetch the complete documentation index at: https://private-7c7dfe99-parallel-read-in-order-multi-part.mintlify.site/llms.txt
> Use this file to discover all available pages before exploring further.

> Таблицы с движком Distributed не хранят собственные данные, но позволяют выполнять распределённую обработку запросов на нескольких серверах. Чтение автоматически распараллеливается. При чтении используются индексы таблиц на удалённых серверах, если они есть.

# Движок таблицы Distributed

<Warning title="Движок Distributed в Cloud">
  Чтобы создать таблицу с движком Distributed в ClickHouse Cloud, можно использовать [табличные функции `remote` и `remoteSecure`](/ru/reference/functions/table-functions/remote).
  Синтаксис `Distributed(...)` нельзя использовать в ClickHouse Cloud.
</Warning>

Таблицы с движком Distributed не хранят собственные данные, но позволяют выполнять распределённую обработку запросов на нескольких серверах.
Чтение автоматически распараллеливается. При чтении используются индексы таблиц на удалённых серверах, если они есть.

## Создание таблицы

```sql theme={null}
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, ...]
```

### Из таблицы

Если distributed таблица `Distributed` указывает на таблицу на текущем сервере, вы можете использовать её схему:

```sql theme={null}
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` и `RemoteSecure` — это постоянные движки таблиц, которые используют те же выражения адресов и учетные данные, что и табличные функции [`remote` and `remoteSecure`](/ru/reference/functions/table-functions/remote):

```sql theme={null}
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`; поддерживаемые сигнатуры см. в их описании. Структуру таблицы можно не указывать — в этом случае она автоматически определяется по удалённой таблице.

[Настройки созданного хранилища](#distributed-settings), такие как `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)`, поскольку табличной функции больше негде его указать; движки эту форму не принимают.

Например:

```sql theme={null}
CREATE TABLE remote_one ENGINE = Remote('127.0.0.1', system, one);
SELECT * FROM remote_one;
```

Это постоянный эквивалент `CREATE TABLE ... AS remote(...)`. Как и табличная функция `remote`, эти движки удобны, но не позволяют декларативно настраивать сегменты и реплики так, как [`Distributed`](#distributed-creating-a-table) в настроенном кластере, поэтому для постоянного, часто используемого набора серверов лучше определить кластер и использовать движок `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

| Параметр | Описание |
| - | - |
| `cluster` | Имя кластера в конфигурационном файле сервера |
| `database` | Имя удаленной базы данных |
| `table` | Имя удаленной таблицы |
| `sharding_key` (Необязательно) | Ключ сегментирования. <br /> Указание `sharding_key` необходимо в следующих случаях: <ul><li>Для операций `INSERT` в таблицу `Distributed` (поскольку движку таблицы нужен `sharding_key`, чтобы определить, как распределить данные). Однако если включена настройка `insert_distributed_one_random_shard`, то для `INSERT` ключ сегментирования не требуется.</li><li>Для использования с `optimize_skip_unused_shards`, поскольку `sharding_key` нужен, чтобы определить, какие сегменты следует запрашивать</li></ul> |
| `policy_name` (Необязательно) | Имя политики; оно будет использоваться для хранения временных файлов при фоновой отправке |

**См. также**

* настройка [distributed\_foreground\_insert](/ru/reference/settings/session-settings/distributed#distributed_foreground_insert)
* [MergeTree](/ru/reference/engines/table-engines/mergetree-family/mergetree#table_engine-mergetree-multiple-volumes) с примерами

### Настройки Distributed

| Настройка | Описание | Значение по умолчанию |
| - | - | - |
| `fsync_after_insert` | Выполнять `fsync` для данных файла после фоновой вставки в Distributed. Гарантирует, что ОС сбросила на диск все вставленные данные **на узле-инициаторе**. | `false` |
| `fsync_directories` | Выполнять `fsync` для каталогов. Гарантирует, что ОС обновила метаданные каталогов после операций, связанных с фоновыми вставками в таблицу `Distributed` (после вставки, после отправки данных в сегмент и т. д.). | `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`](/ru/reference/settings/session-settings/distributed-background#distributed_background_insert_batch) | `0` |
| `background_insert_split_batch_on_failure` | То же, что и [`distributed_background_insert_split_batch_on_failure`](/ru/reference/settings/session-settings/distributed-background#distributed_background_insert_split_batch_on_failure) | `0` |
| `background_insert_sleep_time_ms` | То же, что и [`distributed_background_insert_sleep_time_ms`](/ru/reference/settings/session-settings/distributed-background#distributed_background_insert_sleep_time_ms) | `0` |
| `background_insert_max_sleep_time_ms` | То же, что и [`distributed_background_insert_max_sleep_time_ms`](/ru/reference/settings/session-settings/distributed-background#distributed_background_insert_max_sleep_time_ms) | `0` |
| `flush_on_detach` | Сбрасывать данные на удалённые узлы при `DETACH`/`DROP`/остановке сервера. | `true` |

<Note>
  **Настройки надёжности хранения** (`fsync_...`):

  * Влияют только на фоновые `INSERT` (то есть `distributed_foreground_insert=false`), когда данные сначала сохраняются на диске узла-инициатора, а затем в фоне отправляются в сегменты.
  * Могут значительно снизить производительность `INSERT`
  * Влияют на запись данных, хранящихся в каталоге таблицы `Distributed`, на **узле, который принял вашу вставку**. Если вам нужны гарантии записи данных в базовые таблицы MergeTree, см. настройки надёжности (`...fsync...`) в `system.merge_tree_settings`

  Для **настроек лимитов вставки** (`..._insert`) см. также:

  * настройку [`distributed_foreground_insert`](/ru/reference/settings/session-settings/distributed#distributed_foreground_insert)
  * настройку [`prefer_localhost_replica`](/ru/reference/settings/session-settings/prefer#prefer_localhost_replica)
  * `bytes_to_throw_insert` обрабатывается раньше `bytes_to_delay_insert`, поэтому не следует задавать для него значение меньше, чем `bytes_to_delay_insert`
</Note>

**Пример**

```sql theme={null}
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()`.

## Кластеры

Кластеры настраиваются в [файле конфигурации сервера](/ru/concepts/features/configuration/server-config/configuration-files):

```xml theme={null}
<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`, состоящий из двух сегментов, каждый из которых содержит две реплики. Сегменты — это серверы, содержащие разные части данных (чтобы прочитать все данные, необходимо обратиться ко всем сегментам). Реплики — это серверы-дубликаты (чтобы прочитать все данные, можно обратиться к данным на любой из реплик).

Имена кластеров не должны содержать точек.

Для каждого сервера указываются параметры `host`, `port`, а также при необходимости `user`, `password`, `secure`, `compression`, `bind_host`:

| Параметр | Описание | Значение по умолчанию |
| - | - | - |
| `host` | Адрес удалённого сервера. Можно использовать либо доменное имя, либо IPv4-адрес или IPv6-адрес. Если указан домен, сервер при запуске выполняет DNS-запрос, и результат сохраняется, пока сервер работает. Если DNS-запрос завершается ошибкой, сервер не запускается. Если вы изменили DNS-запись, перезапустите сервер. | - |
| `port` | TCP-порт для обмена сообщениями (`tcp_port` в конфигурации, обычно равен 9000). Не следует путать с `http_port`. | - |
| `user` | Имя пользователя для подключения к удалённому серверу. У этого пользователя должны быть права доступа для подключения к указанному серверу. Доступ настраивается в файле `users.xml`. Дополнительные сведения см. в разделе [Права доступа](/ru/concepts/features/security/access-rights). | `default` |
| `password` | Пароль для подключения к удалённому серверу (не маскируется). | '' |
| `secure` | Следует ли использовать защищённое SSL/TLS‑соединение. Обычно также требуется указать порт (порт по умолчанию для защищённого соединения — `9440`). Сервер должен прослушивать `<tcp_port_secure>9440</tcp_port_secure>` и быть настроен с корректными сертификатами. | `false` |
| `compression` | Использовать сжатие данных. | `true` |
| `bind_host` | Исходный адрес, который следует использовать при подключении к удалённому серверу с этого узла. Поддерживается только IPv4-адрес. Предназначено для сложных сценариев развертывания, когда необходимо задать исходный IP-адрес, используемый ClickHouse для распределённых запросов. | - |

При указании реплик для каждого сегмента при чтении будет выбрана одна из доступных реплик. Вы можете настроить алгоритм балансировки нагрузки (предпочтение, к какой реплике обращаться) — см. настройку [load\_balancing](/ru/reference/settings/session-settings/load-balancing#load_balancing). Если соединение с сервером не удаётся установить, будет предпринята попытка подключения с коротким тайм-аутом. Если подключиться не удалось, будет выбрана следующая реплика, и так для всех реплик. Если попытка подключения не удалась для всех реплик, она будет тем же образом повторена несколько раз. Это повышает устойчивость, но не обеспечивает полной отказоустойчивости: удалённый сервер может принять соединение, но не работать или работать нестабильно.

Вы можете указать только один сегмент (в этом случае обработку запросов следует называть remote, а не distributed) или любое количество сегментов. В каждом сегменте можно указать от одной реплики до любого их числа. Для каждого сегмента можно указать разное количество реплик.

В конфигурации можно указать столько кластеров, сколько потребуется.

Чтобы просмотреть свои кластеры, используйте таблицу `system.clusters`.

Движок `Distributed` позволяет работать с кластером как с локальным сервером. Однако конфигурацию кластера нельзя задавать динамически, её нужно настраивать в конфигурационном файле сервера. Обычно все серверы в кластере имеют одинаковую конфигурацию кластера (хотя это и не обязательно). Кластеры из конфигурационного файла обновляются на лету, без перезапуска сервера.

Если вам нужно каждый раз отправлять запрос неизвестному набору сегментов и реплик, создавать таблицу `Distributed` не нужно — вместо этого используйте табличную функцию `remote`. См. раздел [Табличные функции](/ru/reference/functions/table-functions/index).

## Запись данных

Существует два способа записи данных в кластер:

Во-первых, можно определить, на какие серверы какие данные записывать, и выполнять запись напрямую в каждый сегмент. Иными словами, выполнять прямые операторы `INSERT` в удалённые таблицы кластера, на которые указывает таблица `Distributed`. Это наиболее гибкое решение, поскольку позволяет использовать любую схему шардирования, даже нетривиальную, если этого требует предметная область. Кроме того, это и наиболее оптимальное решение, так как данные можно записывать в разные сегменты полностью независимо друг от друга.

Во-вторых, можно выполнять операторы `INSERT` в таблицу `Distributed`. В этом случае таблица сама распределяет вставленные данные по серверам. Чтобы записывать данные в таблицу `Distributed`, у неё должен быть настроен параметр `sharding_key` (кроме случая, когда сегмент только один).

<Tip>
  Для совместимых запросов `INSERT ... SELECT` между таблицами `Distributed`, использующими один и тот же кластер, [`parallel_distributed_insert_select`](/ru/reference/settings/session-settings/parallel#parallel_distributed_insert_select) может выполнять запрос параллельно на каждом сегменте.
</Tip>

Для каждого сегмента в конфигурационном файле можно определить `<weight>`. По умолчанию вес равен `1`. Данные распределяются по сегментам в объёме, пропорциональном весу сегмента. Все веса сегментов суммируются, затем вес каждого сегмента делится на общую сумму, чтобы определить долю каждого сегмента. Например, если есть два сегмента, и первый имеет вес 1, а второй — вес 2, то в первый будет отправлена одна треть (1 / 3) вставленных строк, а во второй — две трети (2 / 3).

Для каждого сегмента в конфигурационном файле можно определить параметр `internal_replication`. Если этот параметр установлен в `true`, операция записи выбирает первую работоспособную реплику и записывает данные в неё. Используйте это, если таблицы, лежащие в основе таблицы `Distributed`, являются реплицируемыми таблицами (например, используют любой из движков таблиц `Replicated*MergeTree`). Данные будут записаны в одну из реплик таблицы, а затем автоматически реплицированы на остальные реплики.

Если `internal_replication` установлен в `false` (значение по умолчанию), данные записываются во все реплики. В этом случае таблица `Distributed` сама реплицирует данные. Это хуже, чем использование реплицируемых таблиц, поскольку согласованность реплик не проверяется, и со временем они будут содержать немного различающиеся данные.

Чтобы выбрать сегмент, в который будет отправлена строка данных, анализируется выражение шардирования, и берётся остаток от деления на общий вес сегментов. Строка отправляется в тот сегмент, которому соответствует полуинтервал остатков от `prev_weights` до `prev_weights + weight`, где `prev_weights` — это общий вес сегментов с меньшими номерами, а `weight` — вес данного сегмента. Например, если есть два сегмента, и первый имеет вес 9, а второй — вес 10, то строка будет отправлена в первый сегмент для остатков из диапазона \[0, 9), а во второй — для остатков из диапазона \[9, 19).

Выражением шардирования может быть любое выражение из констант и столбцов таблицы, возвращающее целое число. Например, можно использовать выражение `rand()` для случайного распределения данных или `UserID` для распределения по остатку от деления идентификатора пользователя (тогда данные одного пользователя будут находиться на одном сегменте, что упрощает выполнение `IN` и `JOIN` по пользователям). Если один из столбцов распределён недостаточно равномерно, его можно обернуть в хеш-функцию, например `intHash64(UserID)`.

Простой остаток от деления — ограниченное решение для шардирования, и оно подходит не всегда. Оно работает для средних и больших объёмов данных (десятки серверов), но не для очень больших объёмов данных (сотни серверов и более). В последнем случае используйте схему шардирования, подходящую для предметной области, а не таблицы `Distributed`.

О схеме шардирования следует задуматься в следующих случаях:

* Используются запросы, требующие соединения данных (`IN` или `JOIN`) по определённому ключу. Если данные сегментированы по этому ключу, можно использовать локальные `IN` или `JOIN` вместо `GLOBAL IN` или `GLOBAL JOIN`, что значительно эффективнее.
* Используется большое количество серверов (сотни и более) с большим количеством небольших запросов, например запросов по данным отдельных клиентов (сайтов, рекламодателей или партнёров). Чтобы небольшие запросы не затрагивали весь кластер, имеет смысл размещать данные одного клиента на одном сегменте. Как вариант, можно настроить двухуровневое сегментирование: разделить весь кластер на «слои», где слой может состоять из нескольких сегментов. Данные одного клиента располагаются на одном слое, но при необходимости в слой можно добавлять сегменты, и данные распределяются внутри них случайным образом. Для каждого слоя создаются таблицы `Distributed`, а для глобальных запросов создаётся одна общая distributed таблица.

Данные записываются в фоновом режиме. При вставке в таблицу блок данных просто записывается в локальную файловую систему. Затем данные как можно скорее отправляются на удалённые серверы в фоновом режиме. Периодичность отправки данных задаётся настройками [distributed\_background\_insert\_sleep\_time\_ms](/ru/reference/settings/session-settings/distributed-background#distributed_background_insert_sleep_time_ms) и [distributed\_background\_insert\_max\_sleep\_time\_ms](/ru/reference/settings/session-settings/distributed-background#distributed_background_insert_max_sleep_time_ms). Движок `Distributed` отправляет каждый файл со вставленными данными отдельно, но вы можете включить батч-отправку файлов с помощью настройки [distributed\_background\_insert\_batch](/ru/reference/settings/session-settings/distributed-background#distributed_background_insert_batch). Эта настройка повышает производительность кластера за счёт более эффективного использования ресурсов локального сервера и сети. Следует проверять, что данные отправляются успешно, просматривая список файлов (данных, ожидающих отправки) в каталоге таблицы: `/var/lib/clickhouse/data/database/table/`. Количество потоков, выполняющих фоновые задачи, можно задать с помощью настройки [background\_distributed\_schedule\_pool\_size](/ru/reference/settings/server-settings/settings/background#background_distributed_schedule_pool_size).

Если после `INSERT` в таблицу `Distributed` сервер вышел из строя или был аварийно перезапущен (например, из-за аппаратного сбоя), вставленные данные могут быть потеряны. Если в каталоге таблицы обнаружена повреждённая часть данных, она переносится в подкаталог `broken` и больше не используется.

## Чтение данных

При выполнении запроса к таблице `Distributed` запросы `SELECT` отправляются во все сегменты и работают независимо от того, как данные распределены между сегментами (они могут быть распределены совершенно случайным образом). При добавлении нового сегмента переносить в него старые данные не требуется. Вместо этого можно записывать в него новые данные, задав больший вес: данные будут распределены немного неравномерно, но запросы продолжат работать корректно и эффективно.

Когда включена опция `max_parallel_replicas`, обработка запроса распараллеливается между всеми репликами в пределах одного сегмента. Дополнительные сведения см. в разделе [max\_parallel\_replicas](/ru/reference/settings/session-settings/max#max_parallel_replicas).

Чтобы узнать больше о том, как обрабатываются распределенные запросы `in` и `global in`, см. [эту документацию](/ru/reference/statements/in#distributed-subqueries).

## Виртуальные столбцы

#### \_Shard\_num

`_shard_num` — содержит значение `shard_num` из таблицы `system.clusters`. Тип: [UInt32](/ru/reference/data-types/int-uint).

<Note>
  Поскольку табличные функции [`remote`](/ru/reference/functions/table-functions/remote) и [`cluster`](/ru/reference/functions/table-functions/cluster) внутренне создают временную таблицу `Distributed`, `_shard_num` также доступен и в них.
</Note>

**См. также**

* [Описание виртуальных столбцов](/ru/reference/engines/table-engines/index#table_engines-virtual_columns)
* настройка [`background_distributed_schedule_pool_size`](/ru/reference/settings/server-settings/settings/background#background_distributed_schedule_pool_size)
* функции [`shardNum()`](/ru/reference/functions/regular-functions/other-functions#shardNum) и [`shardCount()`](/ru/reference/functions/regular-functions/other-functions#shardCount)
