ClickHouse 1.1.54390 в продакшне: шардирование на нескольких узлах и суб-секундные запросы по миллиардам строк
Разворачиваем ClickHouse для ритейл-аналитики: шардирование на нескольких узлах, репликация через ZooKeeper и первые измерения производительности на реальных данных.
ClickHouse 1.1.54390 - улучшения шардирования и репликации: стабилизация distributed-запросов, починка edge-кейсов в ReplicatedMergeTree
Несколько месяцев назад клиент из ритейла пришёл с задачей, которая на PostgreSQL начинала давать не те тайминги: аналитические запросы по транзакционной истории - сколько каких товаров продано, в каком регионе, в какой час - на объёмах в несколько сотен миллионов строк тянулись от десяти секунд до минуты. BI-система ждёт, аналитики нервничают, регулярные отчёты укладываются в расписание только со скрипом.
Мы начали смотреть в сторону ClickHouse. Внутренний стенд показал интересные результаты, и в марте мы приступили к полноценному развёртыванию на продакшн-данных. Сейчас апрель, кластер работает, и есть что рассказать.
Почему ClickHouse, а не что-то другое
Вариантов рассматривали несколько. Первый - Redshift или BigQuery - сразу отпал: клиент не готов к облаку по внутренним ограничениям. Второй - Greenplum - история длинная, лицензионная и тяжёлая для такого объёма. Третий - Vertica - дорого для пилота.
ClickHouse на этом фоне выглядел разумно: open source, работает на обычных серверах, Яндекс в продакшне держит петабайты. На стенде мы гоняли типичные аналитические запросы клиента по выгрузке из PostgreSQL - агрегации с GROUP BY по нескольким измерениям на 300+ миллионах строк. Результаты были в диапазоне от доли секунды до нескольких секунд, что принципиально отличалось от исходных таймингов.
Как выглядит кластер
Три физических сервера: каждый с несколькими быстрыми дисками, памяти с запасом. Схема - два шарда, каждый шард реплицирован на два узла. ZooKeeper для координации репликации - отдельный трёхузловой кворум на тех же машинах (для старта достаточно, потом вынесем).
Конфигурация кластера в config.xml - примерно так:
<remote_servers>
<retail_cluster>
<shard>
<replica>
<host>ch-node-01</host>
<port>9000</port>
</replica>
<replica>
<host>ch-node-02</host>
<port>9000</port>
</replica>
</shard>
<shard>
<replica>
<host>ch-node-03</host>
<port>9000</port>
</replica>
<replica>
<host>ch-node-04</host>
<port>9000</port>
</replica>
</shard>
</retail_cluster>
</remote_servers>
Таблицы с данными - ReplicatedMergeTree на каждом узле, поверх - Distributed-таблица как единая точка входа для запросов. INSERT-ы идут через Distributed, она сама раскидывает строки по шардам по хеш-ключу.
Схема данных и загрузка
Основная таблица - транзакции продаж. Ключевые поля: дата, магазин, SKU, количество, сумма. Сортировочный ключ (ORDER BY) выбрали по (sale_date, shop_id, sku_id) - это покрывает большинство реальных запросов клиента.
Исторические данные загружали через clickhouse-client с форматом CSV: порциями по несколько миллионов строк, через Distributed-таблицу. На первой загрузке словили интересный эффект - слишком маленькие батчи создавали много мелких кусков (parts), ClickHouse начинал активно их мёрджить и в логах появлялись предупреждения про превышение лимита на количество кусков. Решение: увеличить размер батча до разумных нескольких сотен тысяч строк за раз и дать кластеру время между батчами. Не ракетостроение, но лучше знать заранее.
Итого загрузили около полутора миллиардов строк. На трёх серверах это смотрится скромно - ClickHouse хорошо сжимает колоночные данные, реальный объём на диске в несколько раз меньше исходного CSV.
Что с производительностью
Типичный запрос клиента - продажи по категории за квартал с разбивкой по регионам:
SELECT
region_id,
toStartOfMonth(sale_date) AS month,
sum(amount) AS total_amount,
count() AS tx_count
FROM sales_distributed
WHERE sale_date >= '2017-01-01'
AND sale_date < '2018-01-01'
AND category_id = 42
GROUP BY region_id, month
ORDER BY region_id, month
На PostgreSQL с секционированием - от пятнадцати секунд и выше. На ClickHouse по тому же объёму - меньше секунды, обычно в районе 300-600 мс. Это при том что данные ещё не прогрелись в кешах - ClickHouse сканирует диск.
Более тяжёлые запросы - кросс-продуктовые когорты с несколькими JOIN - занимают несколько секунд. JOIN в ClickHouse требует аккуратности: большие таблицы по обе стороны JOIN работают по-другому, чем в обычных СУБД. Там мы ещё работаем над запросами.
Версия 1.1.54390 и что она принесла
Выпуск 1.1.54390 закрыл несколько неприятных edge-кейсов в шардированных запросах и ReplicatedMergeTree. В частности - стабилизировалось поведение при потере связи с ZooKeeper: узел теперь корректнее переходит в read-only вместо того чтобы продолжать принимать записи и расходиться с репликами. Нас это касается напрямую: при rolling-update ZooKeeper-кворума связь кратковременно рвётся, и раньше надо было следить вручную.
Обновились на неё прямо во время пилота, никаких проблем с апгрейдом не было - rolling restart по одному узлу, кластер не прерывал работу.
Что пока настораживает
ClickHouse - молодой продукт с характером. Несколько вещей держим в голове.
Операционные сложности - ZooKeeper является единой точкой координации репликации. Если ZooKeeper-кворум нездоров, репликация останавливается. Надо мониторить тщательно.
Мутации данных - UPDATE и DELETE в ClickHouse не такие как в PostgreSQL. Исправить строку или удалить несколько записей - это асинхронная операция через ALTER TABLE ... UPDATE/DELETE, которая работает через перезапись кусков. Для аналитического DWH где данные преимущественно append-only это нормально, но клиенту надо понимать это ограничение.
Документация - местами скудная, и часть ответов ищется в исходниках или в чате. Не критично, но закладываем время на это.
Где стоим
Кластер работает в продакшне около трёх недель. BI-система клиента переключена на него для основных аналитических отчётов - аналитики получают ответы там где раньше ждали. На DWH/BI-проектах это первая наша продакшн-инсталляция ClickHouse с шардированием; раньше использовали только в однонодовом режиме для небольших задач.
PostgreSQL у клиента остаётся - как транзакционная база для оперативных данных. ClickHouse берёт на себя аналитическую нагрузку. Разделение разумное и, судя по первым неделям, работает.
Следующий шаг - настроить регулярную подачу свежих данных из PostgreSQL в ClickHouse через ETL, сейчас это делается руками раз в сутки. И нужно определиться с подходом к тем самым тяжёлым JOIN-запросам: либо денормализовывать схему под ClickHouse, либо смотреть на dictionary-таблицы для размерностей.