Цена кластера для пайплайна #13
Notifications
Due Date
No due date set.
Blocks
#14 Кластер: где живёт опыт менти и какая топология
ddmitry/clickstream-ch-kafka-superset-demo
Reference: ddmitry/clickstream-ch-kafka-superset-demo#13
Reference in New Issue
Block a user
Part of #10
Тело не переносилось: issue закрыт задолго до переноса, оригинал остался на GitHub (аккаунт заблокирован 2026-07-26). Название сохранено ради сквозной нумерации и ссылок #NN.
Гист резолюции (из карты): объём средне-крупный (~15–20 файлов, тяжёлое — SQL); главная боль — JOIN поверх Distributed и TRUNCATE в трансформациях (риск для контрольных сумм); приём из Kafka требует одного консьюмера + Distributed-цели; рекомендация — опциональный make up-cluster, не дефолт.
Резолюция (дословный перенос с 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.*_raw—MergeTree(sql/ddl/stg/10_stg.sql:28,41,54,67); плюс 4 таблицыENGINE = Kafkaи 4MATERIALIZED VIEW(там же, строки 77–169).ReplacingMergeTree(src_ingest_ts)(sql/ddl/ods/20_ods.sql:36,79,117,153);4 таблицы ошибок
*_errors—MergeTree(там же, 56,95,132,168).dds.click,dds.event—ReplacingMergeTree(dds_update_ts)(sql/ddl/dds/30_dds.sql:44,76).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:
Локальная (
*_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(иначе базыбудут только на инициаторе).
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 не выполнится.
появятся, если по ошибке задать разные группы — тогда обе ноды прочитают все сообщения дважды.
но MV должен писать в Distributed-цель
stg.*_raw, а не в локальную. Тогда Distributedразложит сырьё по шардам по ключу шардирования, приём остаётся ровно-однократным и повторяемым.
*_raw) не важен для корректности — вниз по потоку ODS/DDS всёравно пересобираются по бизнес-ключам. Можно шардировать по
kafka_offsetили по хешуraw.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). Проблемные места: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там не пишется явно,но модель «полная пересборка» — проверить перед правкой).
(сеть между нодами). Функционально ок, но
now64(3)какdds_update_ts(
30_ods_to_dds.sql:44,135) вычисляется на инициаторе — это хорошо, значение единое.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.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 не сойдётся.То есть корректность контрольных сумм полностью упирается в правильную раскладку/тип джойна.
или задвоится в Kafka-слое,
count/uniqExactразъедутся. Импорт должен оставатьсяexactly-once — отсюда предпочтение «одна нода-консьюмер + Distributed-цель».
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: куда подключать
Сейчас два подключения к одной ноде:
clickhouse://default:123456@clickhouse:9000/default(docker-compose.yml:7, заводитсястроками 181-182) — нативный протокол.
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 файла (по одному на шард).docker-compose.yml:31-32монтированиеz_config.xmlиmacros_ch1.xmlуже закомментировано — место под кластер предусмотрено заранее.
Не переносится / надо дописать самим (в эталоне этого нет):
ON CLUSTERв DDL, пары local+Distributed, правкаTRUNCATE/JOIN(п.1,3).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.sql/ddl/**— вездеON CLUSTER, разбиение на local+Distributed (кроме VIEW).sql/{ods,dds,dm}/*.sql—TRUNCATE ON CLUSTER,GLOBAL/со-локация вджойнах,
GLOBAL NOT IN.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 рисков
GLOBAL/со-локации — тихо неверный результат витрин ирасхождение manifest. Самый коварный: ошибок не бросает, просто цифры не те.
группах; угроза ровно-однократности и, значит, контрольных сумм.
TRUNCATEпо Distributed — полная пересборка ODS/DM сломается, пока не переведут наTRUNCATE … ON CLUSTERпо локальным таблицам.groupArraySample,groupArray) при раскладке связанных строкпо разным шардам — плавающие значения в дашборде между прогонами.
валит
ON CLUSTERDDL и всю пересборку. Плюс рост требований к ресурсам стенда (было 1 нода —станет 2 + Keeper + HAProxy).
make up-cluster(профиль) vs кластер по умолчаниюdocker-compose.ymlиспользовать профили Compose для доп-нод/ZooKeeper/HAProxy; развести DDL итрансформации на два набора (single vs cluster) — либо шаблонизировать SQL (подставлять
ON CLUSTER/Distributedчерез параметр), либо держать вторую ветку файлов. Плюс: учебнаяценность — можно сравнить «одна нода vs кластер» на одном стенде; ничего не ломается для тех, кто
запускает по-старому. Минус: двойное сопровождение SQL, риск расхождения веток, сложнее CI.
каждого запуска, выше порог входа для менти, и любой сбой координатора ломает базовый сценарий.
Для учебного стенда, где ценится быстрый повторяемый прогон, это перебор.
Рекомендация к обсуждению: делать кластер опциональным профилем, а SQL шаблонизировать (один
источник с флагом
ON CLUSTER/local+Distributed), чтобы не плодить две ветки файлов. Это сохраняети простой single-node для менти, и учебную развилку «как оно на кластере».