Распределённый ClickHouse: шардирование, репликация, Keeper

Как ClickHouse масштабируется горизонтально: шардирование и репликация через ClickHouse Keeper, ReplicatedMergeTree и Distributed-таблицы, что ломается при падении узла — реплика недоступна против шард недоступен. Живой кластер 2 шарда × 2 реплики

Пятая статья серии «ClickHouse и аналитические БД». Один узел упирается в свои диск и CPU — дальше начинается распределённый ClickHouse: шардирование (горизонтальное разбиение) и репликация (копии ради надёжности), координируемые через ClickHouse Keeper. Разбираем ReplicatedMergeTree, Distributed-таблицы и что реально происходит при падении узла — на живом кластере из двух шардов по две реплики.

Предыдущие статьи серии разбирали, что происходит внутри одной ноды — MergeTree и гранулы, драйверы, материализованные представления. Здесь нода перестаёт быть единицей масштабирования: данные разбиваются между несколькими нодами (шардирование) и дублируются между нодами внутри одной части разбиения (репликация) — это две независимые оси, и путать их — источник большинства недоразумений про распределённый ClickHouse. Все числа и код ниже — из живого стенда clickhouse/distributed в публичном репозитории digital-cookbook: 4 CH-ноды (2 шарда × 2 реплики) плюс отдельный узел ClickHouse Keeper, собранные отдельным compose-стендом, изолированным от однонодовых стендов остальных статей серии.

Распределённый ClickHouse: 2 шарда × 2 реплики и ClickHouse Keeper, отказ реплики против отказа шарда

В статье

Шардирование vs репликация: разные оси

Проще всего запутаться, если считать шардирование и репликацию вариациями одной идеи «раскидать данные по нескольким машинам». Это не так — у них разные цели, разные механизмы и разная цена отказа:

  • Шардирование — горизонтальное разбиение данных: каждая строка живёт ровно на одном шарде (определяется ключом шардирования), и общий объём данных/нагрузки на запись линейно растёт с числом шардов. Отказ шарда целиком — это отказ доступа к его доле данных, а не ко всем данным сразу, но и не «просто помедленнее»: без явной настройки — ошибка.
  • Репликация — избыточность внутри одного шарда: несколько нод хранят одну и ту же долю данных, чтобы отказ одной ноды не значил отказ доступа к этой доле вообще. Репликация не увеличивает суммарный объём уникальных данных, которые можно хранить кластером, — она увеличивает надёжность и параллелизм чтения ценой лишних копий на диске.

На стенде это два независимых измерения топологии: 2 шарда (масштаб) × 2 реплики на шард (надёжность) = 4 CH-ноды. system.clusters на любой ноде кластера показывает эту матрицу как есть — 4 строки, по паре (shard_num, replica_num) на хост:

SELECT shard_num, replica_num, host_name
FROM system.clusters
WHERE cluster = 'events_cluster'
ORDER BY shard_num, replica_num

На живом прогоне запрос вернул ровно 4 строки: shard=1 replica=1 host=ch-s1-r1, shard=1 replica=2 host=ch-s1-r2, shard=2 replica=1 host=ch-s2-r1, shard=2 replica=2 host=ch-s2-r2 — детерминированная проекция XML-конфига remote_servers.xml, не зависящая от того, есть ли уже какие-либо таблицы в кластере.

flowchart TB subgraph Shard1["Шард 1"] S1R1["ch-s1-r1"] S1R2["ch-s1-r2"] end subgraph Shard2["Шард 2"] S2R1["ch-s2-r1"] S2R2["ch-s2-r2"] end KEEPER["ClickHouse Keeper
координация, кворум"] DIST["Distributed events_distributed
sharding key: cityHash64(user_id)"] DIST -->|fan-out INSERT/SELECT| S1R1 DIST -->|fan-out INSERT/SELECT| S2R1 S1R1 <-->|Keeper-репликация| S1R2 S2R1 <-->|Keeper-репликация| S2R2 S1R1 -.->|координация DDL/репликация| KEEPER S1R2 -.-> KEEPER S2R1 -.-> KEEPER S2R2 -.-> KEEPER style KEEPER fill:#f4d9c6,stroke:#c67a4a style DIST fill:#e8e2d5,stroke:#6d8a99

flowchart TB
  subgraph Shard1["Шард 1"]
    S1R1["ch-s1-r1"]
    S1R2["ch-s1-r2"]
  end
  subgraph Shard2["Шард 2"]
    S2R1["ch-s2-r1"]
    S2R2["ch-s2-r2"]
  end
  KEEPER["ClickHouse Keeper
координация, кворум"] DIST["Distributed events_distributed
sharding key: cityHash64(user_id)"] DIST -->|fan-out INSERT/SELECT| S1R1 DIST -->|fan-out INSERT/SELECT| S2R1 S1R1 <-->|Keeper-репликация| S1R2 S2R1 <-->|Keeper-репликация| S2R2 S1R1 -.->|координация DDL/репликация| KEEPER S1R2 -.-> KEEPER S2R1 -.-> KEEPER S2R2 -.-> KEEPER style KEEPER fill:#f4d9c6,stroke:#c67a4a style DIST fill:#e8e2d5,stroke:#6d8a99
Топология events_cluster: 2 шарда x 2 реплики + ClickHouse Keeper как координатор

ClickHouse Keeper: координация вместо ZooKeeper

Шардирование само по себе не требует внешней координации — каждый шард просто хранит свою долю данных, а решение «какому шарду отправить строку» принимается на лету по ключу шардирования, без сговора между нодами. Репликация — требует: две реплики одного шарда должны согласовывать, какие вставки уже применены, в каком порядке мержить куски, кто выполняет фоновую операцию, а кто её результат подхватывает. Эту координацию исторически делал внешний ZooKeeper; ClickHouse Keeper — собственная реализация того же протокола (ZK wire-протокол, тот же Raft-класс консенсуса), встроенная в сам ClickHouse, без отдельного JVM-процесса рядом.

<clickhouse>
    <zookeeper>
        <node>
            <host>keeper1</host>
            <port>9181</port>
        </node>
    </zookeeper>
</clickhouse>

Тег конфигурации <zookeeper> — исторический, но применим буквально: ClickHouse Keeper говорит по тому же протоколу, поэтому конфигурация подключения не изменилась при переходе с внешнего ZooKeeper на встроенный Keeper — изменилось то, что рядом с CH-нодами теперь поднимается clickhouse-keeper, а не zookeeper, из того же дистрибутива:

<clickhouse>
    <keeper_server>
        <tcp_port>9181</tcp_port>
        <server_id>1</server_id>

        <coordination_settings>
            <operation_timeout_ms>10000</operation_timeout_ms>
            <session_timeout_ms>30000</session_timeout_ms>
        </coordination_settings>

        <raft_configuration>
            <server>
                <id>1</id>
                <hostname>keeper1</hostname>
                <port>9234</port>
            </server>
        </raft_configuration>
    </keeper_server>
</clickhouse>

На стенде Keeper — один узел (raft_configuration с единственным <server>) — осознанное dev-упрощение, а не рекомендация: кворум Raft требует нечётного числа узлов, минимум трёх, чтобы пережить потерю одного без остановки координации. Один узел Keeper здесь достаточен, чтобы показать сам механизм (репликация, дедупликация, DDL ON CLUSTER), но отказ этого единственного узла в стенде не отрабатывается — в проде так разворачивать нельзя, только 3 узла и больше.

Связность с Keeper проверяется до создания каких-либо таблиц — system.zookeeper с путём / уже возвращает служебные znode (keeper, clickhouse), даже без единой ReplicatedMergeTree. Это тот же класс координации, что решает Kafka через собственный встроенный Raft-кворум контроллеров (KRaft) вместо внешнего ZooKeeper — разные протоколы (Keeper говорит по ZK wire-протоколу ради обратной совместимости с существующими клиентами, KRaft — нет), но одна и та же идея: вынести внешнюю координационную систему внутрь самого продукта. Подробнее про эту эволюцию в Kafka, включая живой разбор Raft-кворума на кластере, — в статье «Эксплуатация Kafka: KRaft, мониторинг, тюнинг».

ReplicatedMergeTree и per-нода macros

ReplicatedMergeTree — та же логика частей и фоновых слияний, что у обычного MergeTree, плюс путь в Keeper, через который реплики одного шарда узнают друг о друге. Путь и идентификатор реплики задаются двумя первыми параметрами движка:

CREATE TABLE demo.events ON CLUSTER events_cluster
(
    event_time  DateTime,
    user_id     UInt64,
    event_type  LowCardinality(String),
    url         String,
    duration_ms UInt32,
    country     LowCardinality(String),
    revenue     Decimal(10,2)
)
ENGINE = ReplicatedMergeTree('/clickhouse/tables/{shard}/events', '{replica}')
ORDER BY (event_time, user_id)
PARTITION BY toYYYYMM(event_time)
SETTINGS index_granularity = 8192

Один и тот же текст DDL выполняется на всех 4 нодах через ON CLUSTER events_cluster — фигурные скобки {shard} и {replica} не подставляются клиентом, каждая нода резолвит их сама, локально, из собственного конфига <macros>:

<!-- ch-s1-r1 -->
<clickhouse>
    <macros>
        <shard>1</shard>
        <replica>r1</replica>
        <cluster>events_cluster</cluster>
    </macros>
</clickhouse>

<!-- ch-s1-r2 -->
<clickhouse>
    <macros>
        <shard>1</shard>
        <replica>r2</replica>
        <cluster>events_cluster</cluster>
    </macros>
</clickhouse>

Итог для шарда 1: обе реплики (ch-s1-r1, ch-s1-r2) резолвят путь в /clickhouse/tables/1/events — общий, один и тот же в Keeper, — но с разными значениями {replica} (r1 и r2). Общий путь и есть механизм, которым реплики находят друг друга: любая нода, монтирующая тот же znode-путь, автоматически становится третьей и последующей репликой того же шарда, без изменения DDL. На стенде это подтверждается напрямую через system.zookeeper: SELECT name FROM system.zookeeper WHERE path='/clickhouse/tables/1/events/replicas' после CREATE TABLE возвращает [r1, r2] — Keeper реально знает про обе реплики, а не просто «DDL не упал с ошибкой».

Distributed-таблица и sharding key

ReplicatedMergeTree сам по себе не распределяет строки между шардами — это просто локальная таблица на каждой ноде, реплицируемая внутри своего шарда. Разбиением по шардам занимается отдельный движок — Distributed, который физических данных не хранит вообще, а перенаправляет INSERT/SELECT на нижележащие events по вычисляемому ключу шардирования:

CREATE TABLE demo.events_distributed ON CLUSTER events_cluster
(
    event_time  DateTime,
    user_id     UInt64,
    event_type  LowCardinality(String),
    url         String,
    duration_ms UInt32,
    country     LowCardinality(String),
    revenue     Decimal(10,2)
)
ENGINE = Distributed(events_cluster, demo, events, cityHash64(user_id))

cityHash64(user_id) вычисляется на каждой строке при вставке, и результат хеша определяет, какому шарду она принадлежит — все события одного пользователя всегда попадают на один и тот же шард, что удобно для запросов «по одному пользователю», но не даёт распределить нагрузку какого-то одного «горячего» user_id шире одного шарда. ON CLUSTER создаёт events_distributed на всех 4 нодах разом — любая из них может быть точкой входа: клиент подключается к произвольной ноде, а Distributed сам решает, какой из четырёх реально хранит нужные строки.

Живой прогон вставки 500 000 строк батчами по 100 000 через events_distributed (insert_distributed_sync=1 — синхронная вставка, подробнее ниже) занял 2,512832941 с (198 979 rows/s, характерный прогон, host-зависимо). После сходимости репликации распределение по шардам оказалось shard1=250050, shard2=249950 — близко к равномерному 50/50, что и ожидается от cityHash64: хеш-функция размазывает user_id по диапазону значений равномерно, без привязки к смыслу данных.

insert_distributed_sync: честный недосчёт до сходимости

Здесь стенд наткнулся на явление, которое стоит знать заранее, а не открывать для себя под нагрузкой в проде: insert_distributed_sync=1 гарантирует куда меньше, чем интуитивно ожидается от слова «sync». Сразу после завершения вставки 500 000 строк (insert_distributed_sync=1, батч 100 000) контрольный SELECT shardNum(), count() FROM events_distributed GROUP BY shardNum() показал shard1=250050 shard2=233212 total=483262 — недостача 16 738 строк относительно вставленных 500 000. Это не потеря данных и не гонка в коде стенда — явление воспроизводится детерминированно на этом датасете и seed при каждом прогоне.

Причина в том, что именно гарантирует настройка: insert_distributed_sync=1 заставляет Distributed-таблицу дождаться подтверждения (ACK) от какой-то одной живой реплики каждого шарда на каждый отправляемый батч — какую именно реплику выбрать, решает load_balancing независимо на каждый батч, не обязательно всегда одну и ту же. Синхронность здесь про то, что клиент не получит ответ «вставка прошла», пока эта одна реплика не подтвердит запись, — но она ничего не гарантирует про вторую реплику того же шарда: та узнаёт о новых строках отдельно, через обычную асинхронную Keeper-репликацию, и на это уходит время, пусть и небольшое. На этом прогоне обе пары реплик сошлись быстро — шард 1 и шард 2 конвергировали за 3,9 мс и 3,8 мс соответственно (1 опрос) — но полагаться на то, что fan-out SELECT сразу после вставки увидит согласованную картину, нельзя в принципе, независимо от того, насколько маленькой оказывается задержка на конкретном прогоне.

func pollShardConverged(ctx context.Context, connA, connB clickhouse.Conn, table string, timeout time.Duration) (count uint64, waited time.Duration, polls int) {
    deadline := time.Now().Add(timeout)
    start := time.Now()
    const interval = 200 * time.Millisecond
    for {
        a, _ := countLocal(ctx, connA, table)
        b, _ := countLocal(ctx, connB, table)
        polls++
        if a == b {
            return a, time.Since(start), polls
        }
        if time.Now().After(deadline) {
            log.Fatalf("shard did not converge within %s: A=%d B=%d", timeout, a, b)
        }
        time.Sleep(interval)
    }
}

После сходимости обеих реплик каждого шарда (опрос локальных count() до совпадения, таймаут 30 с, интервал 200 мс — poll-until-converged) картина стабилизируется: shard1=250050 shard2=249950 total=500000 — сумма совпадает с вставленным ровно, повторный fan-out SELECT через Distributed даёт тот же результат. Практический вывод не в том, что insert_distributed_sync бесполезен, — он честно решает свою задачу (гарантия, что данные хотя бы где-то на диске к моменту ответа клиенту), а в том, что «синхронно вставлено» и «мгновенно видно на любой реплике при чтении» — два разных утверждения, и код, которому важна согласованность fan-out чтения сразу после записи, должен закладывать опрос до сходимости или задержку, а не доверять первому же ответу.

Репликация и дедупликация вставок

Отдельно от вставки через Distributed стенд проверяет саму репликацию — напрямую, в обход шардирования: 3000 маркерных строк (фиксированный event_time, не time.Now() — ради байт-в-байт воспроизводимого блока) отправляются одним INSERT напрямую в ch-s1-r1, минуя Distributed-таблицу целиком:

func insertRows(ctx context.Context, conn clickhouse.Conn, table string, rows []row) error {
    insertSQL := fmt.Sprintf("INSERT INTO demo.%s %s", table, insertColumns)
    batch, err := conn.PrepareBatch(ctx, insertSQL)
    if err != nil {
        return fmt.Errorf("prepare batch: %w", err)
    }
    for i, rw := range rows {
        if err := batch.Append(rw.eventTime, rw.userID, rw.eventType, rw.url, rw.durationMs, rw.country, rw.revenue); err != nil {
            return fmt.Errorf("append row %d: %w", i, err)
        }
    }
    return batch.Send()
}

ch-s1-r1 после вставки показывает ожидаемое: count(marker) вырос ровно на 3000 (с 0 до 3000). Интереснее — ch-s1-r2, второй реплике того же шарда, вставка не посылалась вообще, ни напрямую, ни через Distributed. Опрос ch-s1-r2 показал те же 3000/3000 строк — но не мгновенно, а через 212,155406 мс (2 опроса с интервалом 200 мс): именно столько занял путь «ch-s1-r1 записала часть → зарегистрировала в Keeper по общему пути /clickhouse/tables/1/eventsch-s1-r2 увидела новую запись в логе репликации → скачала и применила часть у себя».

Следующий шаг — дедупликация: та же самая, байт-в-байт идентичная пачка из 3000 строк отправляется в ch-s1-r1 вторым, отдельным INSERT. Наивно можно было бы ожидать count(marker) = 6000 — вместо этого он остаётся 3000, без изменений. ReplicatedMergeTree перед тем, как физически записать новую часть, сравнивает хеш вставляемого блока с недавней историей уже применённых блоков (окно контролируется настройкой replicated_deduplication_window, координация — снова через Keeper) и, обнаружив совпадение, отбрасывает блок целиком, как уже применённый. ch-s1-r2 после этого повторного INSERT тоже остаётся на 3000 строках без изменений — реплицировать было нечего: дубликат отброшен на принимающей ноде до того, как он вообще попал в лог репликации.

Важная граница этого механизма: дедупликация по хешу блока — это защита от повторной отправки того же самого блока (например, из-за ретрая клиента после потерянного ACK), а не общая защита от логических дублей на уровне бизнес-данных. Если два разных INSERT содержат разные по составу, но семантически дублирующие друг друга строки (разный порядок, разбиение на другие батчи), ReplicatedMergeTree их дубликатами не сочтёт — это ответственность прикладного кода.

Отказ узла: реплика недоступна vs шард недоступен

Разница между шардированием и репликацией, разобранная в начале статьи, здесь становится измеримой: у отказа одной реплики шарда и отказа шарда целиком — принципиально разная цена для клиента. Стенд проверяет оба случая на baseline shard1=253050 shard2=249950 total=503000 (основной датасет 500 000 строк + 3000 маркерных из предыдущего шага, дедуп не задвоил счёт), останавливая контейнеры конкретных нод командой docker stop снаружи Go-процесса.

sequenceDiagram participant App as Клиент participant Dist as events_distributed (coord ch-s1-r1) participant S1 as Шард 1 (r1 жива, r2 остановлена) participant S2 as Шард 2 (обе реплики остановлены) Note over S1: Сценарий А — одна реплика шарда недоступна App->>Dist: SELECT count() Dist->>S1: fan-out (живая реплика отвечает) Dist->>S2: fan-out (обе реплики живы в этом сценарии) S1-->>Dist: partial count S2-->>Dist: partial count Dist-->>App: total = baseline, без потерь Note over S2: Сценарий Б — весь шард недоступен App->>Dist: SELECT count() (skip_unavailable_shards=0) Dist->>S1: fan-out (ok) Dist--xS2: обе реплики недоступны Dist-->>App: ошибка code 279 (обязателен ответ каждого шарда) App->>Dist: SELECT count() (skip_unavailable_shards=1) Dist->>S1: fan-out (ok) Dist--xS2: обе реплики недоступны, пропущено Dist-->>App: partial = только шард 1, без ошибки

sequenceDiagram
    participant App as Клиент
    participant Dist as events_distributed (coord ch-s1-r1)
    participant S1 as Шард 1 (r1 жива, r2 остановлена)
    participant S2 as Шард 2 (обе реплики остановлены)

    Note over S1: Сценарий А — одна реплика шарда недоступна
    App->>Dist: SELECT count()
    Dist->>S1: fan-out (живая реплика отвечает)
    Dist->>S2: fan-out (обе реплики живы в этом сценарии)
    S1-->>Dist: partial count
    S2-->>Dist: partial count
    Dist-->>App: total = baseline, без потерь

    Note over S2: Сценарий Б — весь шард недоступен
    App->>Dist: SELECT count() (skip_unavailable_shards=0)
    Dist->>S1: fan-out (ok)
    Dist--xS2: обе реплики недоступны
    Dist-->>App: ошибка code 279 (обязателен ответ каждого шарда)

    App->>Dist: SELECT count() (skip_unavailable_shards=1)
    Dist->>S1: fan-out (ok)
    Dist--xS2: обе реплики недоступны, пропущено
    Dist-->>App: partial = только шард 1, без ошибки
Отказ узла: реплика недоступна (без потерь, без настроек) vs шард недоступен (ошибка либо частичный результат, осознанно)

Сценарий А — одна реплика шарда недоступна. docker stop останавливает только ch-s1-r2, вторая реплика того же шарда (ch-s1-r1) остаётся жива. SELECT count() FROM events_distributed после этого возвращает 503000 — ровно baseline, без каких-либо специальных настроек. Distributed прозрачно выбирает живую реплику шарда для fan-out, клиент отказ соседней реплики не замечает вообще — это и есть смысл репликации: цена отказа одной из N реплик — ноль для читателя, пока хотя бы одна реплика шарда жива.

Сценарий Б — весь шард недоступен. docker stop останавливает обе реплики шарда 2 (ch-s2-r1 и ch-s2-r2), шард 1 не трогается. Здесь поведение зависит от явной настройки:

_, errDefault := queryDistributedTotal(boundedCtx, coord) // skip_unavailable_shards не задан (default = 0)
// errDefault != nil: code 279 "All connection tries failed"

skipCtx := clickhouse.Context(ctx, clickhouse.WithSettings(clickhouse.Settings{
    "skip_unavailable_shards": 1,
}))
partial, _ := queryDistributedTotal(skipCtx, coord)
// partial == 253050 (только шард 1, без ошибки)

С дефолтной настройкой (skip_unavailable_shards=0) Distributed обязан получить ответ от каждого шарда — раз обе реплики шарда 2 недостижимы, запрос падает с ошибкой code: 279, All connection tries failed (ATTEMPT_TO_READ_AFTER_EOF в логе) вместо того, чтобы молча вернуть неполные данные. С skip_unavailable_shards=1 тот же запрос возвращает 253050 — ровно данные живого шарда 1, совпадающие с его долей в baseline, и без ошибки, но заведомо меньше полного total=503000. Разница между сценариями А и Б в терминах ответственности: отказ реплики Distributed компенсирует сам, без участия оператора; отказ шарда целиком требует осознанного решения — либо честно упасть (по умолчанию, чтобы неполный агрегат не выдать за полный), либо явно согласиться на частичный ответ через skip_unavailable_shards.

После docker start обеих остановленных нод кластер возвращается к shard1=253050 shard2=249950 total=503000 — полностью совпадает с baseline, без единой довставки: именованные docker-тома пережили stop/start (не down -v), данные никуда не делись, только временно были недоступны по сети.

Что выбрать под задачу

  • Данных больше, чем тянет диск/CPU одной ноды, — нужно шардирование: Distributed поверх нескольких ReplicatedMergeTree, ключ шардирования выбирать по полю, вокруг которого строится большинство запросов (в стенде — user_id, потому что типовой сценарий — «события одного пользователя»).
  • Нужна устойчивость к отказу отдельной ноды без потери доступа к данным — репликация: минимум 2 реплики на шард, координация через Keeper обязательна, дополнительных данных сверх избыточности копий она не добавляет.
  • Координация (DDL ON CLUSTER, репликация, дедупликация вставок) требует Keeper — в проде 3 узла минимум, не 1: единственный узел координатора — единая точка отказа для всего кластера, даже если сами CH-ноды продублированы.
  • Вставка через Distributed с insert_distributed_sync=1 — не гарантия мгновенной согласованности между репликами шарда при чтении сразу после записи; если это важно, закладывать опрос до сходимости на стороне клиента.
  • Отказ реплики — прозрачен для читателя без вмешательства; отказ шарда целиком — требует осознанного выбора между skip_unavailable_shards=0 (упасть, не соврать про полноту) и =1 (отдать честный частичный результат).

Дальше в серии — то, что происходит с этими же parts под операционной нагрузкой: мониторинг, backup, тюнинг codec и то, почему точечные UPDATE/DELETE остаются анти-паттерном даже на распределённом кластере — в статье «Эксплуатация и тюнинг ClickHouse».

Версии в прогоне: clickhouse-server/clickhouse-keeper 26.6.1.1193, clickhouse-go v2.47.0.

Обсуждение в Telegram

Присоединиться →

Комментарии