Партиционирование и шардинг
Шардирование делит данные по узлам, когда один узел не тянет объём или поток записи. Собес проверяет главное решение - выбор ключа шардирования: по нему данные раскладываются и по нему же ищутся.
Стержень: ключ определяет, будут ли запросы бить в один шард или расползаться по всем; монотонный ключ и горячий ключ - две классические катастрофы.
// Формулировки: «как выберешь ключ шардирования?», «что такое scatter-gather?», «что делать с горячим ключом?».
Hash, range, consistent hashing
Hash-шардирование раскладывает равномерно, но диапазонный запрос бежит по всем шардам. Range-шардирование держит диапазоны локально, зато последовательные ключи (время, автоинкремент) создают горячий хвост - вся свежая запись валится в последний шард.
Consistent hashing с виртуальными нодами решает эластичность: добавление узла переселяет примерно 1/N данных, а не всё. Это основа кластеров класса Cassandra и DynamoDB.
// Отсюда первое правило: не шардировать по монотонному ключу. Иначе распределённая система по записи ведёт себя как один перегруженный узел.
- shard key
- ключ раскладки и маршрутизации данных
- consistent hashing
- кольцо хешей: минимум переселений при смене узлов
Горячий ключ и запросы без ключа
Горячий ключ - знаменитость или мегатенант - ломает равномерность, кладя свой шард. Лечат солью к ключу, выделением горячих сущностей отдельно или кэшем перед шардом.
Запрос без ключа шардирования превращается в scatter-gather: опрос всех шардов со сборкой результата - дорого и хрупко. Вторичные индексы в шардированном мире либо локальны (спросить всех), либо глобальны (сами шардированы, но дорогая запись).
// Поэтому ключ выбирают под самые частые запросы: если основной доступ по user_id, а шардируешь по order_id, каждый листинг заказов пользователя станет scatter-gather.
- scatter-gather
- запрос ко всем шардам со сборкой результатов
- hot key
- ключ с аномальной нагрузкой, ломающий баланс
Транзакции и решардинг
Кросс-шардовые транзакции дороги (2PC, саги), поэтому дизайн стремится держать транзакционную границу внутри одного шарда - принцип «все данные пользователя на его шарде». Привычных транзакций через все данные больше нет, это надо закладывать в модель.
Решардинг - плановая операция: онлайн-миграция с двойной записью и чтением по карте маршрутизации. «Потом перешардируем» без готового механизма означает даунтайм.
// Число шардов закладывают с запасом через логические шарды поверх физических: двигать логический шард между машинами куда проще, чем менять саму функцию шардирования на живых данных.
- решардинг
- перераскладка данных при изменении числа шардов
- логический шард
- единица данных поверх физической ноды для лёгкого переезда
Как отвечать: «Как выбрать ключ шардирования?»
Отталкиваюсь от паттерна доступа: по чему чаще всего читают и пишут. Хочу, чтобы типичный запрос попадал в один шард, а не расползался по всем, поэтому ключ обычно совпадает с сущностью, вокруг которой крутится нагрузка - например, user_id, если почти всё привязано к пользователю. Дальше проверяю на две катастрофы. Первая - монотонность: шардировать по времени или автоинкременту нельзя, вся свежая запись пойдёт в один шард. Вторая - перекос и горячие ключи: если есть мегатенант, он положит свой шард, и его надо солить или выносить. Учитываю, что запросы без ключа станут scatter-gather, а кросс-шардовые транзакции дороги, поэтому стараюсь держать транзакционную границу внутри шарда. И закладываю логические шарды с запасом, чтобы расти без пересчёта функции шардирования.
Почему это сильный ответ: ключ выведен из паттерна доступа, названы обе катастрофы (монотонность, перекос), учтены scatter-gather и кросс-шардовые транзакции, и заложен запас на рост.
На чём валят
- −Шардировать по монотонному ключу (автоинкремент, время) - вся запись в последний шард.
- −Забыть про запросы без шард-ключа - каждый листинг стал scatter-gather.
- −Кросс-шардовые транзакции «как раньше» - их больше нет, нужен дизайн границ.
- −Число физических шардов = число логических - рост упирается в решардинг всего.
- −Игнорировать мегатенанта при выборе ключа - один клиент кладёт свой шард.
Проверьте себя
Пять вопросов из банка по этой подтеме. Всего их 12, остальные разбираются в тренажёре.
- Что такое «горячая» партиция (hot partition/shard) и чем она опасна?A)Партиция, которая физически перегревает центральный процессор своего сервера из-за старого железаB)Партиция, которая почти не получает запросов и простаивает без нагрузкиC)Резервная партиция, включающаяся при отказе основной партицииD)Партиция с непропорциональной долей нагрузки — узкое место, сводящее шардинг на нет
показать ответ и разбор
+D)Партиция с непропорциональной долей нагрузки — узкое место, сводящее шардинг на нет// разбор: Горячая партиция получает непропорционально много запросов или данных из-за плохого ключа шардирования (например, по дате — вся свежая нагрузка в одну, или по популярному id). Она становится узким местом: остальные узлы простаивают, а один захлёбывается, и весь смысл горизонтального масштаба теряется. Лечится выбором более равномерного ключа.
- Зачем при шардировании применяют consistent hashing?A)Чтобы при добавлении/удалении узла переезжала лишь малая доля ключей, а не всеB)Чтобы запретить перемещение ключей между узлами при изменении состава кластераC)Чтобы каждый ключ хранился сразу на нескольких узлах кластераD)Чтобы ускорить хеш-функцию за счёт использования более коротких ключей
показать ответ и разбор
+A)Чтобы при добавлении/удалении узла переезжала лишь малая доля ключей, а не все// разбор: При наивном hash(key) % N изменение числа узлов N переселяет почти все ключи — массовый решардинг и простой. Consistent hashing размещает узлы и ключи на кольце, поэтому добавление/удаление узла двигает лишь ~1/N ключей (соседний сегмент). Это делает масштабирование кластера дешёвым и постепенным.
- Почему кросс-шардовые джойны и транзакции — дорогая операция?A)Подобные операции в шардированной системе выполняются так же дёшево и быстро, как локальные, ведь узлы соединены быстрой внутренней сетьюB)Данные лежат на разных узлах — нужны сеть и координация, растут задержка и сложностьC)Джойн между шардами вообще технически не выходит в этом случаеD)Кросс-шардовые операции быстрее локальных за счёт параллелизма узлов
показать ответ и разбор
+B)Данные лежат на разных узлах — нужны сеть и координация, растут задержка и сложность// разбор: Когда нужные данные разбросаны по шардам, джойн или транзакция требуют пересылки данных по сети и распределённой координации (двухфазный коммит), что резко повышает задержку и сложность и снижает надёжность. Поэтому данные стараются шардировать так, чтобы часто связанные строки лежали вместе (co-location по ключу), избегая кросс-шардовых операций.
- Что такое consistent hashing и какую проблему обычного hash(key) % N он решает?A)Consistent hashing полностью убирает необходимость хешировать ключи, распределяя их по алфавитуB)Он обеспечивает, что при изменении числа узлов вообще ни один ключ никуда не переедетC)Кольцо хешей, где при добавлении/удалении узла переезжает лишь малая доля ключей, а не почти все как при % ND)Consistent hashing нужен исключительно для ускорения самой хеш-функции, а не для ребаланса
показать ответ и разбор
+C)Кольцо хешей, где при добавлении/удалении узла переезжает лишь малая доля ключей, а не почти все как при % N// разбор: При шардировании hash(key) % N смена числа узлов N меняет остаток почти для всех ключей — при добавлении узла переезжает огромная доля данных (массовый ремап, лавина миграций). Consistent hashing раскладывает и узлы, и ключи на кольцо: ключ идёт на ближайший по кольцу узел, и при добавлении/удалении узла переезжают только ключи его соседнего участка (~1/N), остальные на месте. Виртуальные узлы выравнивают нагрузку. Основа масштабируемых хранилищ (Dynamo, Cassandra).
- Почему шардированную по user_id таблицу дорого джойнить с шардированной по order_id?A)Джойн дорог потому, что распределённые базы данных не поддерживают операцию JOIN вообщеB)Дороговизна вызвана размером таблиц, а ключи шардирования на стоимость не влияютC)Проблема в индексах: добавив индекс на оба поля, джойн станет дешёвымD)Разные ключи → данные на разных узлах → сетевой shuffle
показать ответ и разбор
+D)Разные ключи → данные на разных узлах → сетевой shuffle// разбор: Локальный джойн возможен, только когда связанные строки лежат на одном узле — то есть таблицы шардированы по одному ключу (co-location). Если одна разложена по user_id, а другая по order_id, строки одного заказа физически на разных узлах, и системе приходится тасовать данные по сети (shuffle/broadcast), что дорого и медленно. Поэтому ключ шардирования выбирают под частые джойны, денормализуют или реплицируют мелкое измерение на все узлы. Кросс-шардовые джойны — общий бич распределённых хранилищ.
дальше
Теорию прочитали. Навык ставится повторением
В Сеньорчике эта подтема идёт в ежедневных сессиях: движок возвращает её, пока ответы не станут уверенными, и ведёт прогресс отдельно по каждой подтеме. Теория внутри тоже бесплатна, лимит только на количество вопросов в день.