сеньорчикОткрыть в Telegram
← вся теориятеория к собесу · ClickHouse

Джойны и распределённые запросы

Зачем это спрашивают

Вставили строку, тут же сделали 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, остальные разбираются в тренажёре.

  1. #ch_distributed1 / 5
    Что такое таблица движка Distributed?
    A)Обычная локальная таблица, которая физически хранит все данные всего кластера целиком на одном узле
    B)Копия таблицы, дублирующая её данные на все реплики ради надёжности и отказоустойчивости системы
    C)«Зонтик» без своих данных: рассылает запрос по шардам и собирает результат, вставки — по ключу шардирования
    D)Временная таблица в оперативной памяти, автоматически удаляемая сразу после завершения запроса
    показать ответ и разбор
    +C)«Зонтик» без своих данных: рассылает запрос по шардам и собирает результат, вставки — по ключу шардирования

    // разбор: Distributed — движок-прокси: сам данных не хранит, а знает топологию кластера и локальные таблицы на шардах. SELECT она разворачивает в подзапросы к каждому шарду и объединяет ответы; INSERT распределяет строки по шардам согласно ключу шардирования. Пользователь пишет как в одну таблицу, а данные лежат на многих узлах. Надёжность (копии) — отдельная задача ReplicatedMergeTree, не Distributed.

  2. #ch_distributed2 / 5
    Зачем в распределённых запросах ClickHouse нужен GLOBAL IN / GLOBAL JOIN?
    A)Чтобы выполнить подзапрос строго на одном шарде и запретить его отправку на остальные узлы кластера
    B)Чтобы полностью отключить шардирование данных на время выполнения одного конкретного запроса к кластеру
    C)Ключевое слово GLOBAL лишь немного ускоряет запрос, но на корректность его результата никак не влияет
    D)Посчитать правый подзапрос один раз и разослать результат на шарды — иначе каждый шард считает его локально и неверно
    показать ответ и разбор
    +D)Посчитать правый подзапрос один раз и разослать результат на шарды — иначе каждый шард считает его локально и неверно

    // разбор: В распределённом запросе обычный IN/JOIN с подзапросом выполняется на КАЖДОМ шарде локально — значит его правая часть видит только данные своего шарда, что даёт неверный результат или лишний трафик. GLOBAL заставляет инициатора вычислить подзапрос один раз и разослать готовый результат (временную таблицу) на все шарды, где он используется целиком. Плата — передача этого результата по сети, поэтому справа держат компактный набор.

  3. #ch_distributed3 / 5
    Как выбирают ключ шардирования (sharding key) для Distributed-вставок?
    A)По колонке с равномерным распределением (часто hash от id), чтобы данные и нагрузка легли по шардам ровно
    B)По возрастающей временной метке события, чтобы данные шли на шарды по очереди
    C)Ключ шардирования выбирать не нужно: ClickHouse сам раскладывает строки поровну
    D)По колонке самой низкой кардинальности, например по булевому флагу активности пользователя
    показать ответ и разбор
    +A)По колонке с равномерным распределением (часто hash от id), чтобы данные и нагрузка легли по шардам ровно

    // разбор: Ключ шардирования определяет, на какой шард уедет строка. Цель — равномерность: распределить данные и запросную нагрузку по шардам без перекосов, поэтому берут высокоэнтропийную колонку или её хеш (intHash от user_id). Плохой выбор (низкая кардинальность, монотонное время) создаёт горячие шарды: всё пишется/читается с одного узла. Ещё учитывают локальность: часто джойнимые вместе данные лучше держать на одном шарде (co-location).

  4. #ch_distributed4 / 5
    Для чего служит секция SAMPLE в запросе к MergeTree?
    A)Чтобы вернуть точный результат, но отсортированный в случайном порядке при каждом запуске
    B)Быстро прикинуть ответ по детерминированной подвыборке данных (нужен SAMPLE BY в схеме таблицы)
    C)Чтобы вставить в таблицу случайные тестовые данные для проверки её общей работоспособности под нагрузкой
    D)Чтобы скопировать структуру таблицы без переноса её реальных данных
    показать ответ и разбор
    +B)Быстро прикинуть ответ по детерминированной подвыборке данных (нужен SAMPLE BY в схеме таблицы)

    // разбор: SAMPLE позволяет считать по части данных ради скорости: SAMPLE 0.1 берёт примерно десятую долю. Работает, если в таблице задан SAMPLE BY (обычно хеш от id) — тогда выборка детерминирована (тот же запрос даёт ту же подвыборку) и пригодна для оценок с домножением. Результат приближённый — для быстрых прикидок на огромных таблицах, а не для точных цифр. Требует продумать SAMPLE BY на этапе схемы.

  5. #ch_distributed5 / 5
    Данные из Kafka грузятся в ReplacingMergeTree, но отчёты иногда показывают задвоение. Наиболее вероятная причина?
    A)ReplacingMergeTree неисправен и не способен убирать дубликаты, нужен полностью другой движок
    B)Причина исключительно в неверно выбранном ключе шардирования Distributed-таблицы данного кластера
    C)Схлопывание идёт только при слиянии кусков — до него запросы без FINAL/argMax видят дубли
    D)Так и должно быть: ReplacingMergeTree по своей конструкции намеренно удваивает каждую строку
    показать ответ и разбор
    +C)Схлопывание идёт только при слиянии кусков — до него запросы без FINAL/argMax видят дубли

    // разбор: Дедуп в ReplacingMergeTree — следствие фонового слияния, а оно происходит когда-нибудь, не сразу после вставки. Свежезалитые из Kafka дубли (или строки в разных кусках) до мержа сосуществуют, и отчёт без FINAL или без агрегации argMax по версии посчитает их дважды. Лечение: считать устойчиво к дублям (GROUP BY + argMax), точечно FINAL, либо гарантировать идемпотентность вставки. «Дедуп при merge» недетерминирован по времени.

дальше

Теорию прочитали. Навык ставится повторением

В Сеньорчике эта подтема идёт в ежедневных сессиях: движок возвращает её, пока ответы не станут уверенными, и ведёт прогресс отдельно по каждой подтеме. Теория внутри тоже бесплатна, лимит только на количество вопросов в день.