Цена кластера для пайплайна #13

Closed
opened 2026-07-26 21:39:14 +03:00 by ddmitry · 1 comment
Owner

Part of #10

Тело не переносилось: issue закрыт задолго до переноса, оригинал остался на GitHub (аккаунт заблокирован 2026-07-26). Название сохранено ради сквозной нумерации и ссылок #NN.

Гист резолюции (из карты): объём средне-крупный (~15–20 файлов, тяжёлое — SQL); главная боль — JOIN поверх Distributed и TRUNCATE в трансформациях (риск для контрольных сумм); приём из Kafka требует одного консьюмера + Distributed-цели; рекомендация — опциональный make up-cluster, не дефолт.

Part of #10 Тело не переносилось: issue закрыт задолго до переноса, оригинал остался на GitHub (аккаунт заблокирован 2026-07-26). Название сохранено ради сквозной нумерации и ссылок #NN. Гист резолюции (из карты): объём средне-крупный (~15–20 файлов, тяжёлое — SQL); главная боль — JOIN поверх Distributed и TRUNCATE в трансформациях (риск для контрольных сумм); приём из Kafka требует одного консьюмера + Distributed-цели; рекомендация — опциональный make up-cluster, не дефолт.
ddmitry added the wayfinder:research label 2026-07-26 21:39:15 +03:00
Author
Owner

Резолюция (дословный перенос с GitHub, автор dementev-dev):

Цена перехода на кластер ClickHouse 2×1 (тикет #13)

Оценка без правок кода. Гипотеза: 2 шарда, реплик нет, плюс координатор (ZooKeeper/ClickHouse Keeper).
Опора: текущий стенд (одна нода ClickHouse) и эталон ~/sources/clickhouse-learning-cluster (пример c2sh2rep).

Короткий вывод: переезд затрагивает почти весь SQL-слой и запуск. Главная боль — не DDL, а JOIN
поверх распределённых таблиц (Distributed) и TRUNCATE в трансформациях. Именно от них зависит,
совпадут ли контрольные суммы manifest после импорта.


1. DDL: движки таблиц, что станет Distributed, где нужен ON CLUSTER, зависимость от координатора

Сейчас всё живёт на одной ноде, координатор не нужен. Инвентарь движков:

  • STG: stg.*_rawMergeTree (sql/ddl/stg/10_stg.sql:28,41,54,67); плюс 4 таблицы
    ENGINE = Kafka и 4 MATERIALIZED VIEW (там же, строки 77–169).
  • ODS: 4 основные — ReplacingMergeTree(src_ingest_ts) (sql/ddl/ods/20_ods.sql:36,79,117,153);
    4 таблицы ошибок *_errorsMergeTree (там же, 56,95,132,168).
  • DDS: dds.click, dds.eventReplacingMergeTree(dds_update_ts) (sql/ddl/dds/30_dds.sql:44,76).
  • DM: витрины — обычные VIEW (sql/ddl/dm/40_dm.sql:26,72,99,115,130,157); плюс одна
    таблица dm.dq_summaryMergeTree (sql/dm/40_dds_to_dm.sql:35).

Что меняется на кластере 2×1:

  • Каждую хранимую таблицу нужно разложить на пару «локальная + Distributed».
    Локальная (*_local) создаётся ON CLUSTER, поверх неё — таблица Distributed(<cluster>, <db>, <local>, <ключ>).
    Пример из эталона: ~/sources/clickhouse-learning-cluster/README.md:99-113.
  • Про движок локальных таблиц. Реплик нет, поэтому по данным ReplicatedMergeTree не обязателен —
    формально хватило бы MergeTree/ReplacingMergeTree как сейчас. Но переход на Replicated*
    всё равно желателен: он нужен, если позже добавят реплику, и он аккуратно работает с ON CLUSTER.
    Эталон использует именно ReplicatedMergeTree даже в примерах (README.md:104).
  • ON CLUSTER нужен всем CREATE/DROP/TRUNCATE, чтобы объект появился/очистился на обеих нодах.
  • Зависимость от координатора. ON CLUSTER использует распределённую очередь DDL, а она живёт в
    ZooKeeper/Keeper. То есть координатор становится обязательным даже без реплик — просто ради того,
    чтобы DDL и TRUNCATE ON CLUSTER расходились по нодам. ReplicatedMergeTree тоже требует Keeper.
  • sql/ddl/00_databases.sql:10-13CREATE DATABASE тоже надо делать ON CLUSTER (иначе базы
    будут только на инициаторе).
  • Скрипт применения DDL scripts/apply_clickhouse_ddl.sh шлёт файлы в один сервис clickhouse
    через docker compose exec (строки 48–68). При ON CLUSTER это ок (один инициатор разложит по
    кластеру), но список нод, кластер и макросы шардов надо где-то определить (см. п.6).

Замечание: витрины DM — это VIEW поверх DDS. VIEW данные не хранит, шардировать его не нужно;
но он должен ссылаться на Distributed-таблицы DDS, иначе увидит только локальный кусок.


2. Kafka-приём на кластере: кто читает, риск дублей, ключ шардирования

Сейчас: 1 брокер, у каждого топика 1 партиция, replication-factor 1
(scripts/load_kafka_data.sh:107-108). Kafka-таблицы и MV — по одному экземпляру на единственной ноде.

Что ломается и какие развилки:

  • Если создать ENGINE = Kafka на обеих нодах с тем же kafka_group_name
    (10_stg.sql:82,93,102,113), брокер поделит партиции между консьюмерами группы. При 1 партиции
    на топик читать будет только одна нода, вторая простаивает. Значит все сырые данные осядут на
    одном шарде — предположение «данные легли на 2 шарда» на уровне STG не выполнится.
  • Риск дублей. Одинаковая группа на двух нодах дублей не даёт (партиции не пересекаются). Дубли
    появятся, если по ошибке задать разные группы — тогда обе ноды прочитают все сообщения дважды.
  • Рекомендуемый для демо и детерминизма вариант: оставить Kafka-таблицу и MV на одной ноде,
    но MV должен писать в Distributed-цель stg.*_raw, а не в локальную. Тогда Distributed
    разложит сырьё по шардам по ключу шардирования, приём остаётся ровно-однократным и повторяемым.
  • Ключ шардирования для сырья (*_raw) не важен для корректности — вниз по потоку ODS/DDS всё
    равно пересобираются по бизнес-ключам. Можно шардировать по kafka_offset или по хешу raw.
  • Ключ шардирования для событий (важно ниже, п.3–4): для dds.event и dds.clickclick_id,
    не event_id. Причина: витрина dm.v_events_enriched джойнит event к click по click_id
    (sql/ddl/dm/40_dm.sql:64). Если оба шардированы по click_id, связанные строки лягут на один
    шард и JOIN станет локальным. event_id как ключ развёл бы событие и его клик по разным шардам.
    Нюанс: event.click_idNullable; события с NULL click_id схлопнутся на один шард (перекос),
    но у них клика и так нет — на корректность не влияет.

Вариант с 2 партициями на топик и чтением обеими нодами возможен, но делает раскладку сырья по
шардам зависимой от того, как брокер назначил партиции консьюмерам — это недетерминированно между
прогонами. Для учебного стенда одна нода-консьюмер + Distributed-цель проще и воспроизводимее.


3. Слои ETL: где заложена одна нода, что нацелить на Distributed

Трансформации — batch-SQL, который Airflow гонит одним ClickHouseOperator на одну ноду
(airflow/dags/etl_pipeline_dag.py). Проблемные места:

  • TRUNCATE перед пересборкой. sql/ods/20_stg_to_ods.sql:26-33 (8 таблиц) и
    sql/dm/40_dds_to_dm.sql:40. На кластере TRUNCATE по Distributed-таблице не проходит — нужно
    TRUNCATE TABLE <db>.<local> ON CLUSTER <cluster>, то есть чистить именно локальные таблицы на
    всех нодах. Это правка каждого TRUNCATE. То же для полной перезагрузки DDS (30_ods_to_dds.sql
    тоже стратегия TRUNCATE+INSERT, см. комментарий строки 18-19; сейчас TRUNCATE там не пишется явно,
    но модель «полная пересборка» — проверить перед правкой).
  • INSERT … SELECT. Вставка в Distributed-таблицу перераспределяет строки по ключу шардирования
    (сеть между нодами). Функционально ок, но now64(3) как dds_update_ts
    (30_ods_to_dds.sql:44,135) вычисляется на инициаторе — это хорошо, значение единое.
  • JOIN поверх Distributed — главный риск. Джойны есть в:
    • 30_ods_to_dds.sql: device ⋈ geo по click_id (строки 74–104), browser ⋈ location по
      event_id (строки 155–172);
    • sql/ddl/dm/40_dm.sql:63-64: dds.event ⋈ dds.click по click_id.
      По умолчанию LEFT JOIN двух Distributed-таблиц без GLOBAL даёт неверный/неполный результат:
      каждый шард джойнит свою левую часть только с локальной правой. Правильно одно из двух:
      (а) GLOBAL LEFT JOIN — правую таблицу собирают на инициаторе и рассылают (просто, но дорого по
      памяти/сети); либо (б) со-локация: шардировать обе стороны по ключу джойна, тогда локальный
      JOIN корректен. Для event⋈click со-локация по click_id (см. п.2). Для browser⋈location
      (event_id) и device⋈geo (click_id) при сборке DDS — либо GLOBAL, либо со-локация ODS-таблиц
      по соответствующим ключам. Это разные ключи для разных пар — со-локацию под все джойны сразу не
      подобрать, часть придётся делать через GLOBAL.
  • Группировки/argMax (30_ods_to_dds.sql, GROUP BY + argMax) поверх Distributed считаются
    корректно: частичная агрегация на шардах, слияние на инициаторе. Дедупликация ReplacingMergeTree
    здесь и так сделана явно через argMax/GROUP BY, а не через FINAL, поэтому раскладка по шардам
    на неё не влияет.
  • dm.dq_summary (40_dds_to_dm.sql) — набор count()/countIf по слоям через UNION ALL.
    Все агрегаты по Distributed считаются верно. Единственное — NOT IN (SELECT …) для «осиротевших»
    событий (строки 105-108) поверх Distributed требует GLOBAL NOT IN, иначе подзапрос выполнится
    локально на каждом шарде и даст ложные «сироты».

4. Детерминизм: сойдутся ли контрольные суммы manifest после импорта

Проверки бьют по одному месту — витрине dm.v_events_enriched:

  • scripts/check_startup_history_manifest.sh:48-60count(), uniqExact(click_id),
    uniqExact(user_domain_id), min/max(event_ts).
  • generator/src/clickstream_generator/airflow_control.py:68-80 (CLICKHOUSE_STATS_SQL) — тот же набор.

Разбор:

  • Сами агрегаты кластер-совместимы: count, uniqExact, min, max над Distributed сливаются из
    частичных состояний шардов и дают тот же результат независимо от раскладки данных. Здесь риска нет.
  • Риск не в агрегатах, а в источнике. Оба запроса читают dm.v_events_enriched — а это тот
    самый LEFT JOIN event⋈click по click_id (п.3). Если джойн на кластере сделан неверно
    (без GLOBAL и без со-локации), витрина недосчитает/переврёт строки, и manifest не сойдётся.
    То есть корректность контрольных сумм полностью упирается в правильную раскладку/тип джойна.
  • Полнота и ровно-однократность приёма (п.2). Если при переезде на кластер сырьё частично потеряется
    или задвоится в Kafka-слое, count/uniqExact разъедутся. Импорт должен оставаться
    exactly-once — отсюда предпочтение «одна нода-консьюмер + Distributed-цель».
  • Недетерминированные места в DM: groupArraySample(1, 1919) в dm.v_session_overview
    (sql/ddl/dm/40_dm.sql:145-146) и groupArray в строках 139–141. Фиксированный seed делает выборку
    повторяемой только если группа (click_id) целиком на одном шарде — тогда состояние считается
    на одной ноде. Если событие и его клик разъехались по шардам, порядок слияния состояний между
    шардами может «плавать» и сэмпл станет недетерминированным. Ещё один довод за со-локацию по
    click_id. На manifest-проверку это прямо не влияет (она эту витрину не читает), но важно для
    стабильности дашборда и прочих сверок.
  • Явных ORDER BY … LIMIT в витринах нет (кроме ORDER BY ключа таблицы dq_summary), так что
    классического «LIMIT без устойчивой сортировки» риска в DM нет.

5. Superset: куда подключать

Сейчас два подключения к одной ноде:

  • Airflow: clickhouse://default:123456@clickhouse:9000/default (docker-compose.yml:7, заводится
    строками 181-182) — нативный протокол.
  • Superset: clickhousedb://default:123456@clickhouse:8123/default
    (superset/init_superset.py:46-79, и захардкожено в superset/dashboards/ecommerce_analytics.zip.json:314).

На кластере обоим клиентам всё равно можно ходить в любую ноду — Distributed-таблица прозрачна с
любого узла. Но правильнее подключать через балансировщик (HAProxy), как в эталоне
(~/sources/clickhouse-learning-cluster/configs/haproxy/haproxy.cfg, порт 8124 HTTP / 9001 native),
чтобы не зависеть от одной ноды. Правки: имя хоста clickhouse → сервис-балансировщик и, возможно,
порт. Отдельно — захардкоженный URI в JSON дашборда (строка 314); хост зашит, при смене адреса его
надо менять руками (или параметризовать при импорте в init_superset.py).


6. Что переносится из эталона clickhouse-learning-cluster, а что нет

Эталон — чистый ClickHouse-кластер (4 ноды, ZooKeeper, HAProxy), Kafka и ETL там нет.

Переносится напрямую (как образец):

  • Структура docker-compose.yml: несколько сервисов clickN, общий zookeeper, haproxy,
    общая сеть, ulimits (у нас ulimits уже такие же).
  • configs/z_config.xml — блок <remote_servers> с кластером c2sh2rep (2 шарда × 2 реплики).
    Для нашей гипотезы 2×1 нужно взять из него только по одной реплике на шард (см. c4sh1rep как
    образец «шард = одна реплика», z_config.xml:14-27) и оставить два таких шарда. Плюс блок
    <zookeeper> (строки 42-47).
  • configs/macros_chN.xml — макросы {shard}/{replica} на каждую ноду
    (macros_ch1.xml); нам нужно 2 файла (по одному на шард).
  • HAProxy-конфиг — почти как есть, только 2 бэкенда вместо 4.
  • Показательно: в нашем docker-compose.yml:31-32 монтирование z_config.xml и macros_ch1.xml
    уже закомментировано — место под кластер предусмотрено заранее.

Не переносится / надо дописать самим (в эталоне этого нет):

  • Kafka-брокер, топики, Kafka-таблицы + MV и их раскладка по кластеру (п.2).
  • Весь ETL: ON CLUSTER в DDL, пары local+Distributed, правка TRUNCATE/JOIN (п.1,3).
  • Airflow (подключение, один инициатор DDL), Superset (балансировщик), скрипты
    apply_clickhouse_ddl.sh, load_kafka_data.sh, check_startup_history_manifest.sh.
  • Эталонные примеры используют реплики (internal_replication=true); у нас реплик нет — берём
    топологию шардов, но не логику реплик.

7. Вердикт: объём, топ-5 рисков, профиль vs дефолт

Грубый объём изменений

Средне-крупно. Затронутые области:

  • Инфраструктура: docker-compose.yml (+2 ноды ClickHouse, ZooKeeper/Keeper, HAProxy), новые
    configs/z_config.xml (2×1), 2× configs/macros_*.xml, configs/haproxy/haproxy.cfg.
  • DDL: 5 файлов sql/ddl/** — везде ON CLUSTER, разбиение на local+Distributed (кроме VIEW).
  • Трансформации: 3 файла sql/{ods,dds,dm}/*.sqlTRUNCATE ON CLUSTER, GLOBAL/со-локация в
    джойнах, GLOBAL NOT IN.
  • Приём: Kafka-таблицы/MV → Distributed-цель; топики (партиции) в scripts/load_kafka_data.sh.
  • Скрипты и проверки: apply_clickhouse_ddl.sh, check_startup_history_manifest.sh,
    подключение Superset (init_superset.py, JSON дашборда), возможно Airflow-connection.
  • Документация: docs/ARCHITECTURE.md, docs/OPERATIONS.md, README.md (по правилу «DDL/инфра →
    доки в том же PR»). Учебная ценность высокая — но и объём пояснений большой.

Ориентир: ~15–20 файлов, из них содержательная переработка SQL — самое трудоёмкое и рискованное.

Топ-5 рисков

  1. JOIN поверх Distributed без GLOBAL/со-локации — тихо неверный результат витрин и
    расхождение manifest. Самый коварный: ошибок не бросает, просто цифры не те.
  2. Приём из Kafka на 2 шардах — либо простой второй ноды (1 партиция), либо дубли при разных
    группах; угроза ровно-однократности и, значит, контрольных сумм.
  3. TRUNCATE по Distributed — полная пересборка ODS/DM сломается, пока не переведут на
    TRUNCATE … ON CLUSTER по локальным таблицам.
  4. Недетерминизм сэмплов/групп (groupArraySample, groupArray) при раскладке связанных строк
    по разным шардам — плавающие значения в дашборде между прогонами.
  5. Зависимость от координатора — новый обязательный сервис (ZooKeeper/Keeper); его недоступность
    валит ON CLUSTER DDL и всю пересборку. Плюс рост требований к ресурсам стенда (было 1 нода —
    станет 2 + Keeper + HAProxy).

make up-cluster (профиль) vs кластер по умолчанию

  • Опциональный профиль (кластер = отдельный режим, single-node остаётся дефолтом): в
    docker-compose.yml использовать профили Compose для доп-нод/ZooKeeper/HAProxy; развести DDL и
    трансформации на два набора (single vs cluster) — либо шаблонизировать SQL (подставлять
    ON CLUSTER/Distributed через параметр), либо держать вторую ветку файлов. Плюс: учебная
    ценность — можно сравнить «одна нода vs кластер» на одном стенде; ничего не ломается для тех, кто
    запускает по-старому. Минус: двойное сопровождение SQL, риск расхождения веток, сложнее CI.
  • Кластер по умолчанию: проще в сопровождении (один набор SQL), но тяжелее по ресурсам для
    каждого запуска, выше порог входа для менти, и любой сбой координатора ломает базовый сценарий.
    Для учебного стенда, где ценится быстрый повторяемый прогон, это перебор.

Рекомендация к обсуждению: делать кластер опциональным профилем, а SQL шаблонизировать (один
источник с флагом ON CLUSTER/local+Distributed), чтобы не плодить две ветки файлов. Это сохраняет
и простой single-node для менти, и учебную развилку «как оно на кластере».

*Резолюция (дословный перенос с GitHub, автор dementev-dev):* # Цена перехода на кластер ClickHouse 2×1 (тикет #13) Оценка без правок кода. Гипотеза: 2 шарда, реплик нет, плюс координатор (ZooKeeper/ClickHouse Keeper). Опора: текущий стенд (одна нода ClickHouse) и эталон `~/sources/clickhouse-learning-cluster` (пример `c2sh2rep`). Короткий вывод: переезд затрагивает почти весь SQL-слой и запуск. Главная боль — не DDL, а `JOIN` поверх распределённых таблиц (`Distributed`) и `TRUNCATE` в трансформациях. Именно от них зависит, совпадут ли контрольные суммы manifest после импорта. --- ## 1. DDL: движки таблиц, что станет Distributed, где нужен ON CLUSTER, зависимость от координатора Сейчас всё живёт на одной ноде, координатор не нужен. Инвентарь движков: - STG: `stg.*_raw` — `MergeTree` (`sql/ddl/stg/10_stg.sql:28,41,54,67`); плюс 4 таблицы `ENGINE = Kafka` и 4 `MATERIALIZED VIEW` (там же, строки 77–169). - ODS: 4 основные — `ReplacingMergeTree(src_ingest_ts)` (`sql/ddl/ods/20_ods.sql:36,79,117,153`); 4 таблицы ошибок `*_errors` — `MergeTree` (там же, 56,95,132,168). - DDS: `dds.click`, `dds.event` — `ReplacingMergeTree(dds_update_ts)` (`sql/ddl/dds/30_dds.sql:44,76`). - DM: витрины — обычные `VIEW` (`sql/ddl/dm/40_dm.sql:26,72,99,115,130,157`); плюс одна таблица `dm.dq_summary` — `MergeTree` (`sql/dm/40_dds_to_dm.sql:35`). Что меняется на кластере 2×1: - Каждую хранимую таблицу нужно разложить на пару «локальная + Distributed». Локальная (`*_local`) создаётся `ON CLUSTER`, поверх неё — таблица `Distributed(<cluster>, <db>, <local>, <ключ>)`. Пример из эталона: `~/sources/clickhouse-learning-cluster/README.md:99-113`. - Про движок локальных таблиц. Реплик нет, поэтому по данным `ReplicatedMergeTree` не обязателен — формально хватило бы `MergeTree`/`ReplacingMergeTree` как сейчас. Но переход на `Replicated*` всё равно желателен: он нужен, если позже добавят реплику, и он аккуратно работает с `ON CLUSTER`. Эталон использует именно `ReplicatedMergeTree` даже в примерах (`README.md:104`). - `ON CLUSTER` нужен всем `CREATE`/`DROP`/`TRUNCATE`, чтобы объект появился/очистился на обеих нодах. - Зависимость от координатора. `ON CLUSTER` использует распределённую очередь DDL, а она живёт в ZooKeeper/Keeper. То есть координатор становится обязательным даже без реплик — просто ради того, чтобы DDL и `TRUNCATE ON CLUSTER` расходились по нодам. `ReplicatedMergeTree` тоже требует Keeper. - `sql/ddl/00_databases.sql:10-13` — `CREATE DATABASE` тоже надо делать `ON CLUSTER` (иначе базы будут только на инициаторе). - Скрипт применения DDL `scripts/apply_clickhouse_ddl.sh` шлёт файлы в один сервис `clickhouse` через `docker compose exec` (строки 48–68). При `ON CLUSTER` это ок (один инициатор разложит по кластеру), но список нод, кластер и макросы шардов надо где-то определить (см. п.6). Замечание: витрины DM — это `VIEW` поверх DDS. `VIEW` данные не хранит, шардировать его не нужно; но он должен ссылаться на Distributed-таблицы DDS, иначе увидит только локальный кусок. --- ## 2. Kafka-приём на кластере: кто читает, риск дублей, ключ шардирования Сейчас: 1 брокер, у каждого топика **1 партиция**, `replication-factor 1` (`scripts/load_kafka_data.sh:107-108`). Kafka-таблицы и MV — по одному экземпляру на единственной ноде. Что ломается и какие развилки: - Если создать `ENGINE = Kafka` **на обеих** нодах с тем же `kafka_group_name` (`10_stg.sql:82,93,102,113`), брокер поделит партиции между консьюмерами группы. При 1 партиции на топик читать будет **только одна нода**, вторая простаивает. Значит все сырые данные осядут на одном шарде — предположение «данные легли на 2 шарда» на уровне STG не выполнится. - Риск дублей. Одинаковая группа на двух нодах дублей не даёт (партиции не пересекаются). Дубли появятся, если по ошибке задать **разные** группы — тогда обе ноды прочитают все сообщения дважды. - Рекомендуемый для демо и детерминизма вариант: оставить Kafka-таблицу и MV **на одной ноде**, но MV должен писать в **Distributed**-цель `stg.*_raw`, а не в локальную. Тогда Distributed разложит сырьё по шардам по ключу шардирования, приём остаётся ровно-однократным и повторяемым. - Ключ шардирования для сырья (`*_raw`) не важен для корректности — вниз по потоку ODS/DDS всё равно пересобираются по бизнес-ключам. Можно шардировать по `kafka_offset` или по хешу `raw`. - Ключ шардирования для событий (важно ниже, п.3–4): для `dds.event` и `dds.click` — **`click_id`**, не `event_id`. Причина: витрина `dm.v_events_enriched` джойнит `event` к `click` по `click_id` (`sql/ddl/dm/40_dm.sql:64`). Если оба шардированы по `click_id`, связанные строки лягут на один шард и `JOIN` станет локальным. `event_id` как ключ развёл бы событие и его клик по разным шардам. Нюанс: `event.click_id` — `Nullable`; события с `NULL click_id` схлопнутся на один шард (перекос), но у них клика и так нет — на корректность не влияет. Вариант с 2 партициями на топик и чтением обеими нодами возможен, но делает раскладку сырья по шардам зависимой от того, как брокер назначил партиции консьюмерам — это недетерминированно между прогонами. Для учебного стенда одна нода-консьюмер + Distributed-цель проще и воспроизводимее. --- ## 3. Слои ETL: где заложена одна нода, что нацелить на Distributed Трансформации — batch-SQL, который Airflow гонит одним `ClickHouseOperator` на одну ноду (`airflow/dags/etl_pipeline_dag.py`). Проблемные места: - **TRUNCATE перед пересборкой.** `sql/ods/20_stg_to_ods.sql:26-33` (8 таблиц) и `sql/dm/40_dds_to_dm.sql:40`. На кластере `TRUNCATE` по Distributed-таблице не проходит — нужно `TRUNCATE TABLE <db>.<local> ON CLUSTER <cluster>`, то есть чистить именно локальные таблицы на всех нодах. Это правка каждого `TRUNCATE`. То же для полной перезагрузки DDS (`30_ods_to_dds.sql` тоже стратегия TRUNCATE+INSERT, см. комментарий строки 18-19; сейчас `TRUNCATE` там не пишется явно, но модель «полная пересборка» — проверить перед правкой). - **INSERT … SELECT.** Вставка в Distributed-таблицу перераспределяет строки по ключу шардирования (сеть между нодами). Функционально ок, но `now64(3)` как `dds_update_ts` (`30_ods_to_dds.sql:44,135`) вычисляется на инициаторе — это хорошо, значение единое. - **JOIN поверх Distributed — главный риск.** Джойны есть в: - `30_ods_to_dds.sql`: `device ⋈ geo` по `click_id` (строки 74–104), `browser ⋈ location` по `event_id` (строки 155–172); - `sql/ddl/dm/40_dm.sql:63-64`: `dds.event ⋈ dds.click` по `click_id`. По умолчанию `LEFT JOIN` двух Distributed-таблиц без `GLOBAL` даёт неверный/неполный результат: каждый шард джойнит свою левую часть только с локальной правой. Правильно одно из двух: (а) `GLOBAL LEFT JOIN` — правую таблицу собирают на инициаторе и рассылают (просто, но дорого по памяти/сети); либо (б) **со-локация**: шардировать обе стороны по ключу джойна, тогда локальный `JOIN` корректен. Для `event⋈click` со-локация по `click_id` (см. п.2). Для `browser⋈location` (`event_id`) и `device⋈geo` (`click_id`) при сборке DDS — либо `GLOBAL`, либо со-локация ODS-таблиц по соответствующим ключам. Это разные ключи для разных пар — со-локацию под все джойны сразу не подобрать, часть придётся делать через `GLOBAL`. - **Группировки/argMax** (`30_ods_to_dds.sql`, `GROUP BY` + `argMax`) поверх Distributed считаются корректно: частичная агрегация на шардах, слияние на инициаторе. Дедупликация `ReplacingMergeTree` здесь и так сделана явно через `argMax`/`GROUP BY`, а не через `FINAL`, поэтому раскладка по шардам на неё не влияет. - `dm.dq_summary` (`40_dds_to_dm.sql`) — набор `count()`/`countIf` по слоям через `UNION ALL`. Все агрегаты по Distributed считаются верно. Единственное — `NOT IN (SELECT …)` для «осиротевших» событий (строки 105-108) поверх Distributed требует `GLOBAL NOT IN`, иначе подзапрос выполнится локально на каждом шарде и даст ложные «сироты». --- ## 4. Детерминизм: сойдутся ли контрольные суммы manifest после импорта Проверки бьют по одному месту — витрине `dm.v_events_enriched`: - `scripts/check_startup_history_manifest.sh:48-60` — `count()`, `uniqExact(click_id)`, `uniqExact(user_domain_id)`, `min/max(event_ts)`. - `generator/src/clickstream_generator/airflow_control.py:68-80` (`CLICKHOUSE_STATS_SQL`) — тот же набор. Разбор: - Сами агрегаты кластер-совместимы: `count`, `uniqExact`, `min`, `max` над Distributed сливаются из частичных состояний шардов и дают тот же результат независимо от раскладки данных. Здесь риска нет. - **Риск не в агрегатах, а в источнике.** Оба запроса читают `dm.v_events_enriched` — а это тот самый `LEFT JOIN event⋈click` по `click_id` (п.3). Если джойн на кластере сделан неверно (без `GLOBAL` и без со-локации), витрина недосчитает/переврёт строки, и manifest **не сойдётся**. То есть корректность контрольных сумм полностью упирается в правильную раскладку/тип джойна. - Полнота и ровно-однократность приёма (п.2). Если при переезде на кластер сырьё частично потеряется или задвоится в Kafka-слое, `count`/`uniqExact` разъедутся. Импорт должен оставаться exactly-once — отсюда предпочтение «одна нода-консьюмер + Distributed-цель». - Недетерминированные места в DM: `groupArraySample(1, 1919)` в `dm.v_session_overview` (`sql/ddl/dm/40_dm.sql:145-146`) и `groupArray` в строках 139–141. Фиксированный seed делает выборку повторяемой **только если группа (`click_id`) целиком на одном шарде** — тогда состояние считается на одной ноде. Если событие и его клик разъехались по шардам, порядок слияния состояний между шардами может «плавать» и сэмпл станет недетерминированным. Ещё один довод за со-локацию по `click_id`. На manifest-проверку это прямо не влияет (она эту витрину не читает), но важно для стабильности дашборда и прочих сверок. - Явных `ORDER BY … LIMIT` в витринах нет (кроме `ORDER BY` ключа таблицы `dq_summary`), так что классического «LIMIT без устойчивой сортировки» риска в DM нет. --- ## 5. Superset: куда подключать Сейчас два подключения к одной ноде: - Airflow: `clickhouse://default:123456@clickhouse:9000/default` (`docker-compose.yml:7`, заводится строками 181-182) — нативный протокол. - Superset: `clickhousedb://default:123456@clickhouse:8123/default` (`superset/init_superset.py:46-79`, и захардкожено в `superset/dashboards/ecommerce_analytics.zip.json:314`). На кластере обоим клиентам всё равно можно ходить в любую ноду — Distributed-таблица прозрачна с любого узла. Но правильнее подключать через балансировщик (HAProxy), как в эталоне (`~/sources/clickhouse-learning-cluster/configs/haproxy/haproxy.cfg`, порт 8124 HTTP / 9001 native), чтобы не зависеть от одной ноды. Правки: имя хоста `clickhouse` → сервис-балансировщик и, возможно, порт. Отдельно — захардкоженный URI в JSON дашборда (строка 314); хост зашит, при смене адреса его надо менять руками (или параметризовать при импорте в `init_superset.py`). --- ## 6. Что переносится из эталона `clickhouse-learning-cluster`, а что нет Эталон — чистый ClickHouse-кластер (4 ноды, ZooKeeper, HAProxy), Kafka и ETL там нет. Переносится напрямую (как образец): - Структура `docker-compose.yml`: несколько сервисов `clickN`, общий `zookeeper`, `haproxy`, общая сеть, `ulimits` (у нас `ulimits` уже такие же). - `configs/z_config.xml` — блок `<remote_servers>` с кластером `c2sh2rep` (2 шарда × 2 реплики). Для нашей гипотезы 2×1 нужно взять из него только по одной реплике на шард (см. `c4sh1rep` как образец «шард = одна реплика», `z_config.xml:14-27`) и оставить два таких шарда. Плюс блок `<zookeeper>` (строки 42-47). - `configs/macros_chN.xml` — макросы `{shard}`/`{replica}` на каждую ноду (`macros_ch1.xml`); нам нужно 2 файла (по одному на шард). - HAProxy-конфиг — почти как есть, только 2 бэкенда вместо 4. - Показательно: в нашем `docker-compose.yml:31-32` монтирование `z_config.xml` и `macros_ch1.xml` **уже закомментировано** — место под кластер предусмотрено заранее. Не переносится / надо дописать самим (в эталоне этого нет): - Kafka-брокер, топики, Kafka-таблицы + MV и их раскладка по кластеру (п.2). - Весь ETL: `ON CLUSTER` в DDL, пары local+Distributed, правка `TRUNCATE`/`JOIN` (п.1,3). - Airflow (подключение, один инициатор DDL), Superset (балансировщик), скрипты `apply_clickhouse_ddl.sh`, `load_kafka_data.sh`, `check_startup_history_manifest.sh`. - Эталонные примеры используют реплики (`internal_replication=true`); у нас реплик нет — берём топологию шардов, но не логику реплик. --- ## 7. Вердикт: объём, топ-5 рисков, профиль vs дефолт ### Грубый объём изменений Средне-крупно. Затронутые области: - Инфраструктура: `docker-compose.yml` (+2 ноды ClickHouse, ZooKeeper/Keeper, HAProxy), новые `configs/z_config.xml` (2×1), 2× `configs/macros_*.xml`, `configs/haproxy/haproxy.cfg`. - DDL: 5 файлов `sql/ddl/**` — везде `ON CLUSTER`, разбиение на local+Distributed (кроме VIEW). - Трансформации: 3 файла `sql/{ods,dds,dm}/*.sql` — `TRUNCATE ON CLUSTER`, `GLOBAL`/со-локация в джойнах, `GLOBAL NOT IN`. - Приём: Kafka-таблицы/MV → Distributed-цель; топики (партиции) в `scripts/load_kafka_data.sh`. - Скрипты и проверки: `apply_clickhouse_ddl.sh`, `check_startup_history_manifest.sh`, подключение Superset (`init_superset.py`, JSON дашборда), возможно Airflow-connection. - Документация: `docs/ARCHITECTURE.md`, `docs/OPERATIONS.md`, `README.md` (по правилу «DDL/инфра → доки в том же PR»). Учебная ценность высокая — но и объём пояснений большой. Ориентир: ~15–20 файлов, из них содержательная переработка SQL — самое трудоёмкое и рискованное. ### Топ-5 рисков 1. **JOIN поверх Distributed без `GLOBAL`/со-локации** — тихо неверный результат витрин и расхождение manifest. Самый коварный: ошибок не бросает, просто цифры не те. 2. **Приём из Kafka на 2 шардах** — либо простой второй ноды (1 партиция), либо дубли при разных группах; угроза ровно-однократности и, значит, контрольных сумм. 3. **`TRUNCATE` по Distributed** — полная пересборка ODS/DM сломается, пока не переведут на `TRUNCATE … ON CLUSTER` по локальным таблицам. 4. **Недетерминизм сэмплов/групп** (`groupArraySample`, `groupArray`) при раскладке связанных строк по разным шардам — плавающие значения в дашборде между прогонами. 5. **Зависимость от координатора** — новый обязательный сервис (ZooKeeper/Keeper); его недоступность валит `ON CLUSTER` DDL и всю пересборку. Плюс рост требований к ресурсам стенда (было 1 нода — станет 2 + Keeper + HAProxy). ### `make up-cluster` (профиль) vs кластер по умолчанию - **Опциональный профиль** (кластер = отдельный режим, single-node остаётся дефолтом): в `docker-compose.yml` использовать профили Compose для доп-нод/ZooKeeper/HAProxy; развести DDL и трансформации на два набора (single vs cluster) — либо шаблонизировать SQL (подставлять `ON CLUSTER`/`Distributed` через параметр), либо держать вторую ветку файлов. Плюс: учебная ценность — можно сравнить «одна нода vs кластер» на одном стенде; ничего не ломается для тех, кто запускает по-старому. Минус: двойное сопровождение SQL, риск расхождения веток, сложнее CI. - **Кластер по умолчанию**: проще в сопровождении (один набор SQL), но тяжелее по ресурсам для каждого запуска, выше порог входа для менти, и любой сбой координатора ломает базовый сценарий. Для учебного стенда, где ценится быстрый повторяемый прогон, это перебор. Рекомендация к обсуждению: делать кластер **опциональным профилем**, а SQL шаблонизировать (один источник с флагом `ON CLUSTER`/local+Distributed), чтобы не плодить две ветки файлов. Это сохраняет и простой single-node для менти, и учебную развилку «как оно на кластере».
Sign in to join this conversation.
1 Participants
Notifications
Due Date
No due date set.
Reference: ddmitry/clickstream-ch-kafka-superset-demo#13