Джойны и распределённые запросы
Вставили строку, тут же сделали SELECT - её нет. Через две секунды есть. Ничего не сломалось: реплики догоняют друг друга асинхронно, и запрос попал на ту, которая ещё не получила кусок. Кластерный ClickHouse полон таких сюрпризов, и все они выводятся из устройства.
Типовые формулировки: «вставили данные, а SELECT их не видит», «почему JOIN на кластере вернул не все строки?», «как выбрать ключ шардирования?».
// Формула кластера: шардирование - про масштаб, репликация - про отказоустойчивость. Distributed - фасад, а не волшебство.
Шарды, реплики и фасад
Шардирование режет данные между узлами: каждый шард держит свою часть, и вместе они дают объём и скорость записи. Репликация копирует каждый шард целиком: это отказоустойчивость и дополнительная мощность на чтение. Кластер из четырёх шардов по две реплики - это восемь узлов, каждый шард хранит четверть данных, каждая четверть существует в двух копиях.
Distributed-таблица данных не хранит. Запрос к ней рассылается на все шарды, каждый считает свою часть, инициатор собирает и доагрегирует результаты. Локальные таблицы на шардах при этом самые обычные MergeTree, и работать можно напрямую с ними.
Вставлять можно двумя способами. Через Distributed - тогда сервер сам раскладывает строки по шардам по ключу шардирования, но делает это асинхронно, через очередь на диске. Напрямую в локальные таблицы - надёжнее и предсказуемее, распределение при этом на совести того, кто пишет.
// Тяжёлая финальная агрегация упирается в инициатора: если GROUP BY высокой кардинальности, каждый шард пришлёт ему миллионы частичных строк, и один узел будет их сливать. Лечится либо правильным шардированием (чтобы группировка была локальной), либо настройкой distributed_group_by_no_merge, когда доагрегация не нужна.
Ключ шардирования: перекос и локальность
Ключ выбирают под то, как будут группировать и джойнить. cityHash64(user_id) кладёт все события одного пользователя на один шард, поэтому GROUP BY user_id выполняется локально на каждом шарде, а инициатор просто складывает готовое.
Ключ по дате - классический антипаттерн. Все сегодняшние события летят на один шард: он принимает сто процентов записи и обслуживает большинство запросов, а остальные три простаивают. Кластер из четырёх узлов работает как один, причём как один перегруженный.
Перекос проверяют до заливки, а не после: посчитать распределение хеша ключа по числу шардов на выборке реальных данных - минутная операция. Если пятая часть трафика приходится на одного крупного клиента, ключ по клиенту создаст горячий шард, который будет держать весь кластер.
// И отдельно: добавление нового шарда не перекладывает старые данные. Они останутся там, где лежали, и новый узел будет получать только свежие записи. Перебалансировка - отдельная спланированная операция, а не побочный эффект расширения.
JOIN на кластере: почему нужен GLOBAL
Обычный JOIN в распределённом запросе выполняется на каждом шарде отдельно, и правую часть каждый шард берёт свою, локальную. Если правая таблица тоже шардирована, каждый узел видит лишь свой кусок справочника, и совпадений находит меньше, чем есть. Результат приходит неполный - и это самое неприятное, потому что запрос не падает и ошибки не пишет.
GLOBAL JOIN и GLOBAL IN решают проблему в лоб: правая часть вычисляется один раз на инициаторе и рассылается всем шардам целиком. Работает всегда, но платит сетью, поэтому правая часть должна быть небольшой.
Второй способ - colocation: шардировать обе таблицы одним и тем же ключом. Тогда все строки, которые могут соединиться, физически лежат на одном узле, и обычный JOIN даёт правильный результат без пересылок. Это решение принимается при проектировании, а не при написании запроса.
// Практическое правило для распределённых запросов: если в подзапросе или джойне участвует ещё одна распределённая таблица, по умолчанию считай результат подозрительным, пока не проверил на маленьких данных, что он совпадает с одношардовым.
- GLOBAL JOIN / GLOBAL IN
- правая часть считается на инициаторе и рассылается всем шардам; без него результат неполный
- colocation
- обе таблицы шардированы одним ключом, поэтому соединение локально на каждом узле
Асинхронность и её цена
Репликация асинхронная: узел принял вставку, записал кусок и ответил клиенту, а копии на других репликах появятся чуть позже. Обычно это доли секунды, под нагрузкой - секунды и минуты. Отчёт, случайно попавший на отставшую реплику, покажет неполные данные и никак об этом не сообщит.
Строгие гарантии покупаются задержкой. insert_quorum заставляет вставку дождаться подтверждения от нескольких реплик, select_sequential_consistency - чтение с достаточно свежей. Оба замедляют работу, поэтому их включают точечно, там, где расхождение действительно дорого.
Мониторинг лага обязателен и смотрится в system.replication_queue: длина очереди и возраст самой старой задачи. Растущая очередь означает, что реплика не догоняет, и дальше будет только хуже.
// Координирует всё это Keeper - кворумный сервис метаданных. Его здоровье равно здоровью записи в реплицированные таблицы: когда кворум Keeper теряется, вставки останавливаются, даже если сами узлы ClickHouse живы и чувствуют себя прекрасно.
- insert_quorum
- вставка подтверждается только после записи на несколько реплик; строгость в обмен на задержку
Как отвечать: «Вставили данные, а SELECT их не видит. Что происходит?»
Главный подозреваемый - асинхронная репликация: читаем с реплики, которая ещё не получила кусок. Проверяю лаг в system.replication_queue, и первым делом выясняю, какой узел и какая реплика вообще ответили на этот SELECT, иначе гадать бессмысленно. Если бизнесу нужна гарантия «читаю своё же», включаю insert_quorum на вставке и select_sequential_consistency на чтении, соглашаясь на рост задержки - точечно, а не для всего трафика. Вторая версия: вставка шла через Distributed, а он раскладывает данные по шардам асинхронно, и строки могут ещё лежать в очереди рассылки; тогда надёжнее писать в локальные таблицы напрямую. Третья, редкая: приложение читает не ту таблицу - фасад с другим набором шардов или локальную вместо распределённой.
Правильный главный подозреваемый, диагностический первый шаг, настройки-гарантии с честно названной ценой и две запасные версии.
На чём валят
- −Джойнить шардированные таблицы без GLOBAL или colocation: результат неполный, а ошибки нет.
- −Шардировать по дате - сегодняшний узел раскалён, остальные простаивают.
- −Читать сразу после вставки с другой реплики и считать это потерей данных.
- −Добавить шард и ждать, что данные выровняются сами: старое останется на месте.
- −Не мониторить лаг репликации - отчёты на отставшей реплике врут молча.
- −Забыть про Keeper: его кворум потерян, и запись встала, хотя сами узлы живы.
Проверьте себя
Пять вопросов из банка по этой подтеме. Всего их 12, остальные разбираются в тренажёре.
- Что такое таблица движка Distributed?A)Обычная локальная таблица, которая физически хранит все данные всего кластера целиком на одном узлеB)Копия таблицы, дублирующая её данные на все реплики ради надёжности и отказоустойчивости системыC)«Зонтик» без своих данных: рассылает запрос по шардам и собирает результат, вставки — по ключу шардированияD)Временная таблица в оперативной памяти, автоматически удаляемая сразу после завершения запроса
показать ответ и разбор
+C)«Зонтик» без своих данных: рассылает запрос по шардам и собирает результат, вставки — по ключу шардирования// разбор: Distributed — движок-прокси: сам данных не хранит, а знает топологию кластера и локальные таблицы на шардах. SELECT она разворачивает в подзапросы к каждому шарду и объединяет ответы; INSERT распределяет строки по шардам согласно ключу шардирования. Пользователь пишет как в одну таблицу, а данные лежат на многих узлах. Надёжность (копии) — отдельная задача ReplicatedMergeTree, не Distributed.
- Зачем в распределённых запросах ClickHouse нужен GLOBAL IN / GLOBAL JOIN?A)Чтобы выполнить подзапрос строго на одном шарде и запретить его отправку на остальные узлы кластераB)Чтобы полностью отключить шардирование данных на время выполнения одного конкретного запроса к кластеруC)Ключевое слово GLOBAL лишь немного ускоряет запрос, но на корректность его результата никак не влияетD)Посчитать правый подзапрос один раз и разослать результат на шарды — иначе каждый шард считает его локально и неверно
показать ответ и разбор
+D)Посчитать правый подзапрос один раз и разослать результат на шарды — иначе каждый шард считает его локально и неверно// разбор: В распределённом запросе обычный IN/JOIN с подзапросом выполняется на КАЖДОМ шарде локально — значит его правая часть видит только данные своего шарда, что даёт неверный результат или лишний трафик. GLOBAL заставляет инициатора вычислить подзапрос один раз и разослать готовый результат (временную таблицу) на все шарды, где он используется целиком. Плата — передача этого результата по сети, поэтому справа держат компактный набор.
- Как выбирают ключ шардирования (sharding key) для Distributed-вставок?A)По колонке с равномерным распределением (часто hash от id), чтобы данные и нагрузка легли по шардам ровноB)По возрастающей временной метке события, чтобы данные шли на шарды по очередиC)Ключ шардирования выбирать не нужно: ClickHouse сам раскладывает строки поровнуD)По колонке самой низкой кардинальности, например по булевому флагу активности пользователя
показать ответ и разбор
+A)По колонке с равномерным распределением (часто hash от id), чтобы данные и нагрузка легли по шардам ровно// разбор: Ключ шардирования определяет, на какой шард уедет строка. Цель — равномерность: распределить данные и запросную нагрузку по шардам без перекосов, поэтому берут высокоэнтропийную колонку или её хеш (intHash от user_id). Плохой выбор (низкая кардинальность, монотонное время) создаёт горячие шарды: всё пишется/читается с одного узла. Ещё учитывают локальность: часто джойнимые вместе данные лучше держать на одном шарде (co-location).
- Для чего служит секция SAMPLE в запросе к MergeTree?A)Чтобы вернуть точный результат, но отсортированный в случайном порядке при каждом запускеB)Быстро прикинуть ответ по детерминированной подвыборке данных (нужен SAMPLE BY в схеме таблицы)C)Чтобы вставить в таблицу случайные тестовые данные для проверки её общей работоспособности под нагрузкойD)Чтобы скопировать структуру таблицы без переноса её реальных данных
показать ответ и разбор
+B)Быстро прикинуть ответ по детерминированной подвыборке данных (нужен SAMPLE BY в схеме таблицы)// разбор: SAMPLE позволяет считать по части данных ради скорости: SAMPLE 0.1 берёт примерно десятую долю. Работает, если в таблице задан SAMPLE BY (обычно хеш от id) — тогда выборка детерминирована (тот же запрос даёт ту же подвыборку) и пригодна для оценок с домножением. Результат приближённый — для быстрых прикидок на огромных таблицах, а не для точных цифр. Требует продумать SAMPLE BY на этапе схемы.
- Данные из Kafka грузятся в ReplacingMergeTree, но отчёты иногда показывают задвоение. Наиболее вероятная причина?A)ReplacingMergeTree неисправен и не способен убирать дубликаты, нужен полностью другой движокB)Причина исключительно в неверно выбранном ключе шардирования Distributed-таблицы данного кластераC)Схлопывание идёт только при слиянии кусков — до него запросы без FINAL/argMax видят дублиD)Так и должно быть: ReplacingMergeTree по своей конструкции намеренно удваивает каждую строку
показать ответ и разбор
+C)Схлопывание идёт только при слиянии кусков — до него запросы без FINAL/argMax видят дубли// разбор: Дедуп в ReplacingMergeTree — следствие фонового слияния, а оно происходит когда-нибудь, не сразу после вставки. Свежезалитые из Kafka дубли (или строки в разных кусках) до мержа сосуществуют, и отчёт без FINAL или без агрегации argMax по версии посчитает их дважды. Лечение: считать устойчиво к дублям (GROUP BY + argMax), точечно FINAL, либо гарантировать идемпотентность вставки. «Дедуп при merge» недетерминирован по времени.
дальше
Теорию прочитали. Навык ставится повторением
В Сеньорчике эта подтема идёт в ежедневных сессиях: движок возвращает её, пока ответы не станут уверенными, и ведёт прогресс отдельно по каждой подтеме. Теория внутри тоже бесплатна, лимит только на количество вопросов в день.