сеньорчикОткрыть в Telegram
← вся теориятеория к собесу · Распределённые системы

Партиционирование и шардинг

Шардирование: ключ решает всё

Шардирование делит данные по узлам, когда один узел не тянет объём или поток записи. Собес проверяет главное решение - выбор ключа шардирования: по нему данные раскладываются и по нему же ищутся.

Стержень: ключ определяет, будут ли запросы бить в один шард или расползаться по всем; монотонный ключ и горячий ключ - две классические катастрофы.

// Формулировки: «как выберешь ключ шардирования?», «что такое 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, остальные разбираются в тренажёре.

  1. #partitioning_sharding1 / 5
    Что такое «горячая» партиция (hot partition/shard) и чем она опасна?
    A)Партиция, которая физически перегревает центральный процессор своего сервера из-за старого железа
    B)Партиция, которая почти не получает запросов и простаивает без нагрузки
    C)Резервная партиция, включающаяся при отказе основной партиции
    D)Партиция с непропорциональной долей нагрузки — узкое место, сводящее шардинг на нет
    показать ответ и разбор
    +D)Партиция с непропорциональной долей нагрузки — узкое место, сводящее шардинг на нет

    // разбор: Горячая партиция получает непропорционально много запросов или данных из-за плохого ключа шардирования (например, по дате — вся свежая нагрузка в одну, или по популярному id). Она становится узким местом: остальные узлы простаивают, а один захлёбывается, и весь смысл горизонтального масштаба теряется. Лечится выбором более равномерного ключа.

  2. #partitioning_sharding2 / 5
    Зачем при шардировании применяют consistent hashing?
    A)Чтобы при добавлении/удалении узла переезжала лишь малая доля ключей, а не все
    B)Чтобы запретить перемещение ключей между узлами при изменении состава кластера
    C)Чтобы каждый ключ хранился сразу на нескольких узлах кластера
    D)Чтобы ускорить хеш-функцию за счёт использования более коротких ключей
    показать ответ и разбор
    +A)Чтобы при добавлении/удалении узла переезжала лишь малая доля ключей, а не все

    // разбор: При наивном hash(key) % N изменение числа узлов N переселяет почти все ключи — массовый решардинг и простой. Consistent hashing размещает узлы и ключи на кольце, поэтому добавление/удаление узла двигает лишь ~1/N ключей (соседний сегмент). Это делает масштабирование кластера дешёвым и постепенным.

  3. #partitioning_sharding3 / 5
    Почему кросс-шардовые джойны и транзакции — дорогая операция?
    A)Подобные операции в шардированной системе выполняются так же дёшево и быстро, как локальные, ведь узлы соединены быстрой внутренней сетью
    B)Данные лежат на разных узлах — нужны сеть и координация, растут задержка и сложность
    C)Джойн между шардами вообще технически не выходит в этом случае
    D)Кросс-шардовые операции быстрее локальных за счёт параллелизма узлов
    показать ответ и разбор
    +B)Данные лежат на разных узлах — нужны сеть и координация, растут задержка и сложность

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

  4. #partitioning_sharding4 / 5
    Что такое consistent hashing и какую проблему обычного hash(key) % N он решает?
    A)Consistent hashing полностью убирает необходимость хешировать ключи, распределяя их по алфавиту
    B)Он обеспечивает, что при изменении числа узлов вообще ни один ключ никуда не переедет
    C)Кольцо хешей, где при добавлении/удалении узла переезжает лишь малая доля ключей, а не почти все как при % N
    D)Consistent hashing нужен исключительно для ускорения самой хеш-функции, а не для ребаланса
    показать ответ и разбор
    +C)Кольцо хешей, где при добавлении/удалении узла переезжает лишь малая доля ключей, а не почти все как при % N

    // разбор: При шардировании hash(key) % N смена числа узлов N меняет остаток почти для всех ключей — при добавлении узла переезжает огромная доля данных (массовый ремап, лавина миграций). Consistent hashing раскладывает и узлы, и ключи на кольцо: ключ идёт на ближайший по кольцу узел, и при добавлении/удалении узла переезжают только ключи его соседнего участка (~1/N), остальные на месте. Виртуальные узлы выравнивают нагрузку. Основа масштабируемых хранилищ (Dynamo, Cassandra).

  5. #partitioning_sharding5 / 5
    Почему шардированную по user_id таблицу дорого джойнить с шардированной по order_id?
    A)Джойн дорог потому, что распределённые базы данных не поддерживают операцию JOIN вообще
    B)Дороговизна вызвана размером таблиц, а ключи шардирования на стоимость не влияют
    C)Проблема в индексах: добавив индекс на оба поля, джойн станет дешёвым
    D)Разные ключи → данные на разных узлах → сетевой shuffle
    показать ответ и разбор
    +D)Разные ключи → данные на разных узлах → сетевой shuffle

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

дальше

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

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