Шардирование с vshard
Шардирование в Tarantool реализовано в модуле vshard. Краткое руководство по vshard см. в разделе
Создание шардированного кластера.
Модуль vshard не входит в основной дистрибутив Tarantool. Чтобы установить модуль, выполните команду:
$ tt rocks install vshard
Если вы разрабатываете приложение для шардированного кластера, добавьте зависимость от модуля vshard в файл *.rockspec:
dependencies = {'vshard == 0.1.27'}
Настройка параметров, связанных с шардированием, может включать следующие шаги:
- Настройте параметры подключения, чтобы экземпляры в шардированном кластере могли взаимодействовать друг с другом.
- Укажите, какую роль выполняет каждый набор реплик в шардированном кластере.
- Настройте способ разделения данных по шардам.
- Укажите параметры, связанные с балансировкой данных.
В этом разделе описаны параметры подключения, которые обеспечивают взаимодействие между экземплярами в шардированном кластере. Общие сведения о подключениях см. в разделе Подключения.
В конфигурации шардированного кластера необходимо указать, как маршрутизатор и балансировщик подключаются к хранилищам, с
помощью параметра
iproto.advertise.sharding.
В приведенном ниже примере для этого используется пользователь storage:
iproto:advertise:peer:login: replicatorsharding:login: storage
Пользователь storage должен иметь роль sharding, описанную в следующем разделе.
Чтобы маршрутизатор и балансировщик могли подключаться к хранилищам, следует использовать пользователя с
ролью sharding. В приведенном ниже
примере показано, как назначить роль sharding пользователю storage:
credentials:users:replicator:password: 'topsecret'roles: [replication]storage:password: 'secret'roles: [sharding]
Роль sharding предоставляет различные привилегии в зависимости от
роли шардирования набора реплик. Для наборов реплик с ролью шардирования
storage роль sharding предоставляет следующие привилегии:
- Все привилегии, предоставляемые ролью
replication. - Выполнение функций vshard.storage.*.
Если набор реплик не имеет роли шардирования storage, роль sharding не предоставляет никаких привилегий.
Каждый набор реплик в шардированном кластере может выполнять одну из трех ролей:
router: набор реплик выполняет роль роутера.storage: набор реплик выполняет роль хранилища.rebalancer: набор реплик выполняет роль балансировщика.
Чтобы назначить определенную роль набору реплик или группе наборов реплик, используйте параметр
sharding.roles.
В приведенном ниже примере все наборы реплик в группе storages имеют роль storage, а наборы реплик в группе
routers – роль router.
groups:storages:sharding:roles: [storage]# ...routers:sharding:roles: [router]# ...
Обратите внимание, что роль rebalancer является необязательной. Если она не указана, балансировщик выбирается
автоматически из числа мастер-экземпляров наборов реплик. Чтобы указать балансировщик вручную или отключить его,
используйте параметр
sharding.rebalancer_mode.
В этом разделе описаны параметры конфигурации, связанные с шардированием данных. О том, как определить спейсы для шардирования, см. в разделе Определение данных.
Чтобы задать общее количество сегментов в кластере, настройте параметр
sharding.bucket_count
на глобальном уровне. В примере ниже для
параметра sharding.bucket_count задано значение 1000:
sharding:bucket_count: 1000
Значение параметра sharding.bucket_count должно быть на несколько порядков больше потенциального количества узлов
кластера с учетом возможного масштабирования в будущем.
Если оценочное количество узлов в кластере равно N, то набор данных следует разделить на 100N или даже 1000N сегментов в зависимости от планируемого масштабирования. Это число превышает потенциальное количество узлов кластера в проектируемой системе.
Следует учитывать, что слишком большое количество сегментов может привести к необходимости выделять больше памяти для хранения информации о маршрутизации. С другой стороны, недостаточное количество сегментов снижает гранулярность при балансировке.
Вес набора реплик определяет емкость хранения набора реплик: чем больше вес, тем больше сегментов может хранить набор реплик. Вес набора реплик настраивается с помощью параметра sharding.weight. Этот параметр позволяет хранить основной объем данных на наборе реплик с большим объемом памяти. Также можно задать нулевой вес для набора реплик, чтобы инициировать миграцию его сегментов на остальные узлы кластера.
В примере ниже набор реплик storage-a может хранить в два раза больше данных, чем storage-b:
# ...replicasets:storage-a:sharding:weight: 2# ...storage-b:sharding:weight: 1# ...
Существует эталонное число сегментов в наборе реплик ("эталонный" в данном случае значит идеальный). Если во всем наборе реплик это число остается неизменным, то сегменты распределяются равномерно.
Эталонное число рассчитывается автоматически с учетом количества сегментов в кластере и весов наборов реплик.
Балансировка запускается, если предел дисбаланса набора реплик превышает предел дисбаланса, заданный в конфигурации (параметр sharding.rebalancer_disbalance_threshold).
Предел дисбаланса набора реплик рассчитывается следующим образом:
|эталонное_число_сегментов - текущее_число_сегментов| / эталонное_число_сегментов * 100
Например, кластер настроен следующим образом:
-
Количество сегментов (параметр sharding.bucket_count) равно 3000.
-
Веса трех наборов реплик равны 1, 0.5 и 1.5. В этом случае эталонные числа сегментов для наборов реплик равны:
- 1-й набор реплик – 1000.
- 2-й набор реплик – 500.
- 3-й набор реплик – 1500.
Чтобы инициировать миграцию сегментов набора реплик на остальные узлы кластера, можно задать значение 0 для веса набора реплик. Чтобы инициировать миграцию сегментов из существующих наборов реплик, можно добавить новый набор реплик с ненулевым весом.
При добавлении нового шарда для миграции сегментов на новый шард необходимо перезагрузить конфигурацию на каждом экземпляре:
-
Если используется централизованное хранилище конфигурации, Tarantool перезагружает измененную конфигурацию автоматически.
-
Если используется локальный файл конфигурации, необходимо сначала перезагрузить конфигурацию на всех роутерах, а затем на всех хранилищах.
Изначально в vshard был довольно простой rebalancer – один процесс на одном узле, который вычислял маршруты отправки
сегментов: сколько и кому. Узлы применяли эти маршруты по очереди, последовательно.
К сожалению, такая простая схема работала недостаточно быстро, особенно для Vinyl'а, где затраты ресурсов на чтение диска были сопоставимы с сетевыми затратами. На самом деле, механизм применения маршрутов в балансировщике Vinyl'а большую часть времени был в режиме ожидания.
Теперь каждый узел может параллельно посылать несколько сегментов по кругу в несколько пунктов назначения или всего в один.
Чтобы задать степень параллельности, используйте параметр sharding.rebalancer_max_sending:
sharding:rebalancer_max_sending: 5
У вас уже есть 10 наборов реплик, добавили новый. Теперь все 10 наборов реплик будут пытаться отправить сегменты на новый.
Предположим, что каждый набор реплик может одновременно отправлять до 5 сегментов. В этом случае новый набор реплик испытает довольно большую нагрузку, так как 50 сегментов будут загружаться одновременно. Если узлу нужно выполнять какую-либо другую работу, такая большая нагрузка может быть нежелательной. Кроме того, слишком большое количество параллельно передаваемых сегментов может привести к тайм-аутам в самом процессе балансировки.
Чтобы решить эту проблему, можно задать меньшее значение параметра rebalancer_max_sending для старых наборов реплик или
уменьшить значение параметра rebalancer_max_receiving для нового. В последнем случае некоторые воркеры на старых узлах
будут приостановлены, что можно будет увидеть в логах.
Параметр rebalancer_max_sending важен, если в кластере есть ограничения на максимальное количество сегментов, которые
могут быть одновременно доступны только для чтения. Как вы помните, при отправке сегмент не принимает новые запросы на
запись.
У вас есть 100 000 сегментов, и каждый сегмент хранит ~ 0,001% ваших данных. В кластере 10 наборов реплик. И нельзя
позволить себе заблокировать для записи > 0,1% данных. Таким образом, не следует устанавливать значение
rebalancer_max_sending > 10 на этих узлах. Тогда балансировщик не будет посылать более 100 сегментов одновременно
по всему кластеру.
Если значение rebalancer_max_sending слишком высокое, а значение rebalancer_max_receiving слишком низкое, некоторые
сегменты будут пытаться переместиться – и им это не удастся. Это приведет к лишним затратам сетевых ресурсов и времени
Важно настроить эти параметры так, чтобы они не конфликтовали друг с другом.
Блокировка набора реплик
(параметр sharding.lock) делает
набор реплик невидимым для rebalancer: заблокированный набор реплик не может ни принимать новые сегменты, ни мигрировать
свои собственные сегменты.
Закрепление сегмента (API метод vshard.storage.bucket_pin(bucket_id)) блокирует миграцию конкретного сегмента: закрепленный сегмент остается в наборе реплик, к которому он закреплен, пока не будет откреплен.
Закрепление всех сегментов в наборе реплик не означает блокирование набора реплик. Даже после закрепления всех сегментов незаблокированный набор реплик может принимать новые сегменты.
Блокировка набора реплик полезна, например, для отделения набора реплик от рабочих наборов реплик в целях тестирования или для сохранения некоторых метаданных приложения, которые не должны быть сегментированы в течение некоторого времени. Закрепление сегмента используется в аналогичных случаях, но в меньшем масштабе.
Заблокировав набор реплик и закрепив все сегменты, можно полностью изолировать набор реплик.
Заблокированные наборы реплик и закрепленные сегменты влияют на алгоритм балансировки, так как балансировщик должен
игнорировать заблокированные наборы реплик и учитывать закрепленные сегменты при попытке достичь наилучшего возможного
баланса.
Это нетривиальная задача, поскольку пользователь может закрепить слишком много сегментов в наборе реплик, так что становится невозможным достижение идеального баланса. Например, рассмотрим следующий кластер (предположим, что все веса наборов реплик равны 1).
Начальная конфигурация:
rs1: bucket_count = 150 -- число сегментовrs2: bucket_count = 150, pinned_count = 120 -- число сегментов, число закрепленных сегментов
Добавление нового набора реплик:
rs1: bucket_count = 150rs2: bucket_count = 150, pinned_count = 120rs3: bucket_count = 0
Идеальным балансом было бы 100 - 100 - 100, чего невозможно достичь, поскольку набор реплик rs2 содержит 120
закрепленных сегментов.
Наилучший возможный баланс здесь следующий:
rs1: bucket_count = 90rs2: bucket_count = 120, pinned_count 120rs3: bucket_count = 90
С помощью rebalancer было перемещено максимально возможное количество сегментов из rs2 для уменьшения дисбаланса. При
этом были учтены равные веса rs1 и rs3.
Алгоритмы реализации блокировки и закрепления совершенно разные, хотя с точки зрения функций они похожи.
Заблокированные наборы реплик не участвуют в балансировке. Это означает, что даже если фактическое общее количество сегментов не равно эталонному, устранить дисбаланс из-за блокировки невозможно. Когда балансировщик обнаруживает, что один из наборов реплик заблокирован, он пересчитывает эталонное количество сегментов для незаблокированных наборов реплик так, как если бы заблокированного набора реплик и его сегментов не существовало вообще.
Балансировка наборов реплик с закрепленными сегментами требует более сложного алгоритма. Здесь pinned_count[o] – это
число закрепленных сегментов, а etalon_count – это эталонное число сегментов для набора реплик:
rebalancerвычисляет эталонное число сегментов так, как если бы ни один сегмент не был закреплен. Затем балансировщик проверяет каждый набор реплик и сравнивает эталонное число сегментов с числом закрепленных сегментов в наборе реплик. Еслиpinned_count < etalon_count, незаблокированные наборы реплик (на этом этапе все заблокированные наборы реплик уже отфильтрованы) с закрепленными сегментами могут получать новые сегменты.- Если
pinned_count > etalon_count, дисбаланс невозможно устранить, так какrebalancerне может переместить закрепленные сегменты из этого набора реплик. В таком случае эталонное число обновляется и приравнивается к числу закрепленных сегментов. Наборы реплик сpinned_count > etalon_countне обрабатываются балансировщиком, а число закрепленных сегментов вычитается из общего числа сегментов. Балансировщик пытается переместить как можно больше сегментов из таких наборов реплик. - Эта процедура повторяется с шага 1 для наборов реплик с
pinned_count >= etalon_countдо тех пор, пока не будет выполнятьсяpinned_count <= etalon_countдля всех наборов реплик. Процедура также перезапускается при изменении общего числа сегментов.
Псевдокод для данного алгоритма будет следующим:
function cluster_calculate_perfect_balance(replicasets, bucket_count)-- балансировка сегментов с учетом весов жизнеспособных наборов реплик --end;cluster = <все незаблокированные наборы реплик>;bucket_count = <общее число сегментов в кластере>;can_reach_balance = falsewhile not can_reach_balance docan_reach_balance = truecluster_calculate_perfect_balance(cluster, bucket_count);foreach replicaset in cluster doif replicaset.perfect_bucket_count <replicaset.pinned_bucket_count thencan_reach_balance = falsebucket_count -= replicaset.pinned_bucket_count;replicaset.perfect_bucket_count =replicaset.pinned_bucket_count;end;end;end;cluster_calculate_perfect_balance(cluster, bucket_count);
Сложность алгоритма составляет O(N^2), где N – количество наборов реплик. На каждом шаге алгоритм либо завершает
вычисление, либо игнорирует хотя бы один новый набор реплик, перегруженный закрепленными сегментами, и обновляет эталонное
число сегментов в других наборах реплик.
Ссылка в сегменте – это счетчик в оперативной памяти, который похож на закрепление сегмента со следующими отличиями:
-
Ссылка на сегмент не является постоянной. Ссылки предназначены для запрета перемещения сегмента во время выполнения запроса, но при перезапуске все запросы сбрасываются.
-
Существует два типа ссылок на сегменты: только для чтения (RO) и для чтения-записи (RW).
Если в сегменте есть ссылки типа RW, его нельзя перемещать. Однако, если балансировщику требуется отправка этого сегмента, он блокирует его для новых запросов на запись, ожидает завершения всех текущих запросов, а затем отправляет сегмент.
Если в сегменте есть ссылки типа RO, его можно отправить, но нельзя удалить. Такой сегмент может даже перейти в статус мусора GARBAGE или отправки SENT, но его данные сохраняются до тех пор, пока не уйдет последний читатель.
В одном сегменте могут быть ссылки как типа RO, так и типа RW.
-
В ссылку на сегмент включён счетчик.
Методы vshard.storage.bucket_ref/unref()
вызываются автоматически при использовании методов
vshard.router.call() или
vshard.storage.call().
При использовании низкоуровневого API вида r = vshard.router.route() r:callro/callrw необходимо явно вызывать метод
bucket_ref() внутри функции. Также следует убедиться, что bucket_unref() вызывается после bucket_ref(), иначе
сегмент не сможет быть перемещен с хранилища до перезапуска экземпляра.
Чтобы узнать количество ссылок в сегменте, используйте метод
vshard.storage.buckets_info([идентификатор_сегмента])
(параметр идентификатор_сегмента необязателен).
Пример:
vshard.storage.buckets_info(1)---- 1:status: activeref_rw: 1ref_ro: 1ro_lock: truerw_lock: trueid: 1
Шардированные спейсы должны быть определены в приложении хранилища внутри функции box.once() и должны иметь поле со значениями идентификатор сегмента. Это поле должно соответствовать следующим требованиям:
- Тип данных поля может быть
unsigned,numberилиinteger. - Поле не может содержать NULL.
- Поле должно быть проиндексировано с помощью индекса
shard_index.
Имя этого индекса по умолчанию –
bucket_id. В примере ниже спейсbandsсодержит полеbucket_id, которое используется для распределения набора данных между разными экземплярами хранилища:
box.once('bands', function()box.schema.create_space('bands', {format = {{ name = 'id', type = 'unsigned' },{ name = 'bucket_id', type = 'unsigned' },{ name = 'band_name', type = 'string' },{ name = 'year', type = 'unsigned' }},if_not_exists = true})box.space.bands:create_index('id', { parts = { 'id' }, if_not_exists = true })box.space.bands:create_index('bucket_id', { parts = { 'bucket_id' }, unique = false, if_not_exists = true })end)
Пример на GitHub: sharded_cluster
Все операции DML с данными должны выполняться через роутер с использованием функций vshard.router.call, таких как
vshard.router.callrw()
или
vshard.router.callro().
Например, в приложении хранилища есть функция insert_band, используемая для вставки новых кортежей:
function insert_band(id, bucket_id, band_name, year)box.space.bands:insert({ id, bucket_id, band_name, year })end
В приложении роутера можно определить функцию put, которая определяет, как роутер выбирает хранилище для записи данных:
function put(id, band_name, year)local bucket_id = vshard.router.bucket_id_mpcrc32({ id })vshard.router.callrw(bucket_id, 'insert_band', { id, bucket_id, band_name, year })end
Подробнее см. в разделе Обработка запросов.
Идемпотентные запросы дают одинаковый результат при каждом выполнении. Например, запрос на чтение данных или умножение на единицу – идемпотентные операции. Соответственно, инкремент на единицу – пример неидемпотентной операции. При повторном применении такой операции значение поля увеличивается на 2, а не на 1.
Запрос может потребоваться выполнить повторно, если на стороне сервера или клиента возникла ошибка. В этом случае:
-
Запросы на чтение могут выполняться повторно. Для этого в методе vshard.router.call() (с
mode=read) используется параметрrequest_timeout(начиная сvshard0.1.28). Параметрыrequest_timeoutиtimeoutнеобходимо передавать вместе, при этом должно соблюдаться следующее условие:timeout > request_timeoutНапример, если
timeout = 10иrequest_timeout = 2, в течение 10 секунд роутер может сделать 5 попыток (по 2 секунды каждая) отправки запроса разным репликам, пока запрос не завершится успешно. -
Запросы на запись (vshard.router.callrw()) обычно не могут быть выполнены повторно без проверки того, что они не были применены ранее. Отсутствие такой проверки может привести к появлению дубликатов записей или непредвиденным изменениям данных.
Например, клиент отправил запрос на сервер и ожидает ответ в течение заданного таймаута. Если сервер отправляет успешный ответ по истечении этого времени, клиент не получит этот ответ из-за таймаута и сочтет запрос неудачным. При повторном выполнении этого запроса без дополнительной проверки операция может быть применена дважды.
Запрос на запись может выполняться повторно без проверки в двух случаях:
-
Запрос идемпотентен.
-
Достоверно известно, что предыдущий запрос завершился ошибкой до выполнения операций записи. Например, сервер вернул ошибку ER_READONLY. В этом случае известно, что запрос не мог быть выполнен из-за того, что сервер находится в режиме только для чтения.
-
Примеры дедупликации
Чтобы обеспечить идемпотентность запросов на запись (INSERT, UPDATE, UPSERT и автоинкремент), следует реализовать проверку того, что запрос применяется впервые.
Например, при добавлении нового кортежа в спейс можно использовать уникальный идентификатор вставки для проверки запроса. В примере ниже в рамках одной транзакции:
- Проверяется, существует ли кортеж с идентификатором
keyв спейсеbands. - Если кортежа с таким идентификатором в спейсе нет, кортеж вставляется.
box.begin()if box.space.bands:get{key} == nil thenbox.space.bands:insert{key, value}endbox.commit()
Для запросов на обновление и upsert можно создать спейс дедупликации, в котором будут сохраняться идентификаторы
запросов. Спейс дедупликации – это пользовательский спейс, содержащий список уникальных идентификаторов. Каждый
идентификатор соответствует одному примененному запросу. Этот спейс может иметь любое имя, в примере он
называется deduplication.
В примере ниже в рамках одной транзакции:
- Проверяется, существует ли идентификатор запроса
deduplication_keyв спейсеdeduplication. - Если такого идентификатора нет, он добавляется в спейс дедупликации.
- Если запрос не был применен ранее, указанное поле в спейсе
bandsувеличивается на единицу. Такой подход гарантирует, что каждый запрос на изменение данных будет выполнен только один раз.
function update_1(deduplication_key, key)box.begin()if box.space.deduplication:get{deduplication_key} == nil thenbox.space.deduplication:insert{deduplication_key}box.space.bands:update(key, {{'+', 'value', 1 }})endbox.commit()end
В случае отказа мастера в наборе реплик рекомендуется:
- Переключите одну из реплик в режим мастера. Это позволит новому мастеру обрабатывать все входящие запросы.
- Обновите конфигурацию всех участников кластера. Это перенаправит все запросы новому мастеру.
В случае сбоя всего набора реплик часть набора данных становится недоступной. При этом маршрутизатор пытается переподключиться к мастеру отказавшего набора реплик. Таким образом, после того как набор реплик снова заработает, кластер восстанавливается автоматически.
Для проведения запланированной остановки мастера в наборе реплик рекомендуется:
- Обновите конфигурацию, чтобы использовать другой экземпляр в качестве мастера.
- Перезагрузите конфигурацию на всех экземплярах. После этого все запросы будут перенаправлены новому мастеру.
- Остановите старый мастер.
Для проведения запланированной остановки набора реплик рекомендуется:
- Перенесите все сегменты на другие хранилища кластера. Это можно сделать, назначив набору реплик нулевой вес, что инициирует перенос его сегментов на оставшиеся узлы кластера.
- Обновите конфигурацию всех узлов.
- Остановите набор реплик.
Поиск сегментов, восстановление сегментов и балансировка сегментов выполняются автоматически и не требуют ручного вмешательства.
С технической точки зрения есть несколько файберов, которые отвечают за различные типы действий:
- Файбер обнаружения на маршрутизаторе выполняет поиск сегментов в фоновом режиме.
- Файбер отказоустойчивости на маршрутизаторе поддерживает соединения с репликами.
- Файбер сборщика мусора на каждом мастер-хранилище удаляет содержимое перемещенных сегментов.
- Файбер восстановления сегментов на каждом мастер-хранилище восстанавливает сегменты в состояниях SENDING и RECEIVING в случае перезагрузки.
- Файбер балансировщика на одном мастер-хранилище среди всех наборов реплик выполняет процесс балансировки. Для получения подробной информации см. разделы Процесс балансировки и Миграция сегментов.
Файбер сборщик мусора работает в фоновом режиме на мастер-хранилищах в каждом наборе реплик. Он начинает удалять
содержимое сегмента в состоянии мусора GARBAGE по частям. Когда сегмент пуст, запись о нем удаляется из системного спейса
_bucket.
Файбер восстановления сегмента работает на мастер-хранилищах. Он помогает восстановить сегменты в статусах отправки SENDING и получения RECEIVING в случае перезагрузки.
Сегменты в статусе SENDING восстанавливаются следующим образом:
- Система сначала выполняет поиск сегментов в состоянии SENDING.
- Если такой сегмент найден, система отправляет запрос в целевой набор реплик.
- Если сегмент на целевом наборе реплик находится в состоянии ACTIVE, исходный сегмент удаляется с узла-источника.
Сегменты в статусе RECEIVING удаляются без дополнительных проверок.
Файбер отказоустойчивости работает на каждом маршрутизаторе. Если мастер набора реплик становится недоступным, файбер отказоустойчивости перенаправляет запросы на чтение к репликам. Запросы на запись отклоняются с ошибкой, пока мастер не станет доступным.