В системе заказов один крупный клиент создаёт почти весь трафик, хотя остальные клиенты распределены равномерно. Как изменить ключ шардирования, чтобы убрать горячий шард?
Нельзя использовать только идентификатор клиента как ключ шардирования: крупный клиент создаёт горячий шард. Добавьте к клиенту равномерно изменяющийся компонент, например номер виртуального бакета, вычисляемый из идентификатора заказа, чтобы записи одного клиента распределялись по нескольким шардам.
Такой подход устраняет перегрузку записи, но усложняет чтение всех заказов клиента: запрос должен обращаться к нескольким шардам, а затем объединять результаты.
Шардирование появилось как способ горизонтально распределять данные и нагрузку между несколькими узлами, когда один узел перестаёт справляться с объёмом хранения, операций или сетевого трафика. Ключ шардирования определяет, на какой узел попадёт запись и где искать её при чтении.
Простой ключ обычно выбирают по наиболее естественной границе данных: клиенту, пользователю или организации. Это упрощает запросы одного владельца, но предполагает, что нагрузка между владельцами примерно равномерна. В многопользовательских системах такое предположение часто нарушается из-за крупных клиентов.
Если все заказы клиента направляются на один шард, размер клиента становится фактическим ограничением масштабирования. Даже при большом количестве остальных клиентов один шард может получать непропорционально много операций записи, занимать больше места и иметь худшие задержки.
Увеличение числа шардов само по себе проблему не решает: маршрутизация по прежнему ключу продолжит отправлять весь поток крупного клиента на тот же узел. Репликация может увеличить пропускную способность чтения, но обычно не устраняет узкое место записи и неравномерное распределение хранения.
Используйте составной логический ключ вида идентификатор клиента плюс номер бакета. Номер бакета получают из идентификатора заказа или другого равномерно распределённого атрибута. Физический шард выбирается хешированием этого составного ключа либо через таблицу маршрутизации виртуальных бакетов.
Например, для клиента можно определить несколько бакетов, а каждый новый заказ направлять в бакет по хешу его идентификатора. Тогда поток одного клиента распределяется между несколькими шардами. Число бакетов должно учитывать пиковую нагрузку, размер заказа, допустимое число обращений при чтении и возможность дальнейшего масштабирования.
Запрос одного заказа остаётся адресным, если в нём известны клиент и идентификатор заказа. Запрос всех заказов клиента становится распределённым: система обращается ко всем его бакетам, выполняет объединение результатов, сортировку и пагинацию. Поэтому для таких запросов нужны курсоры, ограничение глубины страницы и контроль числа параллельных обращений.
Важно сохранять в записи полный ключ маршрутизации или иметь детерминированное правило его вычисления. Если бакеты меняются со временем, понадобится таблица диапазонов или версий маршрутизации; простая смена числа бакетов при хешировании может сделать старые записи недоступными по прежнему адресу.
Есть два основных компромисса. Фиксированное число бакетов проще для чтения и маршрутизации, но может оказаться недостаточным при росте клиента. Динамическое выделение бакетов экономит ресурсы, но требует управления метаданными, миграций и согласованного поведения при изменении маршрутов.
Не следует бездумно добавлять случайный суффикс: он распределяет записи, но лишает систему возможности однозначно определить бакет по идентификатору заказа, если этот суффикс не сохраняется в записи или не вычисляется детерминированно. Также нужно отдельно учитывать вторичные индексы, уникальные ограничения и операции, затрагивающие данные нескольких бакетов.
В условной системе заказов ключом шардирования был идентификатор клиента. После подключения крупного маркетплейса почти все операции записи стали попадать на один шард, хотя суммарный объём данных был распределён по кластеру относительно равномерно.
Рассматривались три варианта. Вертикальное увеличение узла быстро внедряется, но даёт временный предел и не устраняет зависимость от одного шарда. Репликация улучшает чтение, однако не решает перегрузку записи. Полный перенос крупного клиента в отдельный кластер изолирует его нагрузку, но усложняет эксплуатацию и создаёт специальное правило маршрутизации.
Выбрали несколько виртуальных бакетов для крупных клиентов с распределением по идентификатору заказа. Для обычных клиентов сохранили прежнюю схему, чтобы не увеличивать стоимость запросов без необходимости. Чтение списка заказов выполнялось по всем бакетам клиента с ограниченной параллельностью, а пагинация использовала составной курсор.
В результате нагрузка записи перестала концентрироваться на одном шарде. Цена решения — более сложные запросы по истории крупного клиента и необходимость контролировать число бакетов при его дальнейшем росте.
Число бакетов оценивают не только по объёму данных, но прежде всего по пиковому числу операций, допустимой нагрузке одного шарда и стоимости fan-out-чтения. Если один шард должен выдерживать не более определённого потока записи, минимальное число бакетов получают из отношения пикового потока клиента к этой границе, затем добавляют запас.
Слишком мало бакетов оставляет горячие участки. Слишком много увеличивает число обращений при чтении списка заказов, объём метаданных и стоимость перебалансировки. Обычно полезно предусмотреть запас и возможность назначать отдельным клиентам собственное число бакетов.
Уникальность, проверяемая только внутри одного шарда, не гарантирует глобальную уникальность. Если номер заказа генерируется независимо в разных бакетах, возможны коллизии, поэтому идентификатор должен быть глобально уникальным сам по себе либо включать компонент, однозначно связанный с маршрутизацией.
Глобальный последовательный счётчик создаёт отдельную точку координации и может стать узким местом. На практике выбирают идентификаторы, уникальность которых обеспечивается без общего счётчика, либо выносят проверку в специализированный координируемый механизм, принимая его стоимость и задержку.
Обычно вводят версию маршрутизации, начинают записывать новые данные по новой схеме, а чтение временно умеет искать данные в старом и новом расположении. Затем существующие записи переносят пакетами с контролем полноты, повторяемостью операций и ограничением влияния на рабочую нагрузку.
После сверки количества и контрольных выборок старое расположение можно вывести из чтения и удалить поэтапно. Нельзя просто изменить формулу ключа: старые записи останутся на прежних шардах, а запросы по новой формуле перестанут их находить.