Tarantool CE/EE Documentation portal logo
Помощь
Обновлена 15 сентября 2026 г. в 08:55

Шардирование с vshard

Шардирование в Tarantool реализовано в модуле vshard. Краткое руководство по vshard см. в разделе Создание шардированного кластера.

Установка

Модуль vshard не входит в основной дистрибутив Tarantool. Чтобы установить модуль, выполните команду:

$ tt rocks install vshard

Если вы разрабатываете приложение для шардированного кластера, добавьте зависимость от модуля vshard в файл *.rockspec:

dependencies = {    'vshard == 0.1.27'}

Обзор конфигурации

Настройка параметров, связанных с шардированием, может включать следующие шаги:

  1. Настройте параметры подключения, чтобы экземпляры в шардированном кластере могли взаимодействовать друг с другом.
  2. Укажите, какую роль выполняет каждый набор реплик в шардированном кластере.
  3. Настройте способ разделения данных по шардам.
  4. Укажите параметры, связанные с балансировкой данных.

Связность

В этом разделе описаны параметры подключения, которые обеспечивают взаимодействие между экземплярами в шардированном кластере. Общие сведения о подключениях см. в разделе Подключения.

Объявляемый URI

В конфигурации шардированного кластера необходимо указать, как маршрутизатор и балансировщик подключаются к хранилищам, с помощью параметра iproto.advertise.sharding. В приведенном ниже примере для этого используется пользователь storage:

iproto:  advertise:    peer:      login: replicator    sharding:      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 не предоставляет никаких привилегий.

Роли шардирования

Каждый набор реплик в шардированном кластере может выполнять одну из трех ролей:

Чтобы назначить определенную роль набору реплик или группе наборов реплик, используйте параметр 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

Пример 1

У вас уже есть 10 наборов реплик, добавили новый. Теперь все 10 наборов реплик будут пытаться отправить сегменты на новый.

Предположим, что каждый набор реплик может одновременно отправлять до 5 сегментов. В этом случае новый набор реплик испытает довольно большую нагрузку, так как 50 сегментов будут загружаться одновременно. Если узлу нужно выполнять какую-либо другую работу, такая большая нагрузка может быть нежелательной. Кроме того, слишком большое количество параллельно передаваемых сегментов может привести к тайм-аутам в самом процессе балансировки.

Чтобы решить эту проблему, можно задать меньшее значение параметра rebalancer_max_sending для старых наборов реплик или уменьшить значение параметра rebalancer_max_receiving для нового. В последнем случае некоторые воркеры на старых узлах будут приостановлены, что можно будет увидеть в логах.

Параметр rebalancer_max_sending важен, если в кластере есть ограничения на максимальное количество сегментов, которые могут быть одновременно доступны только для чтения. Как вы помните, при отправке сегмент не принимает новые запросы на запись.

Пример 2

У вас есть 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 – это эталонное число сегментов для набора реплик:

  1. rebalancer вычисляет эталонное число сегментов так, как если бы ни один сегмент не был закреплен. Затем балансировщик проверяет каждый набор реплик и сравнивает эталонное число сегментов с числом закрепленных сегментов в наборе реплик. Если pinned_count < etalon_count, незаблокированные наборы реплик (на этом этапе все заблокированные наборы реплик уже отфильтрованы) с закрепленными сегментами могут получать новые сегменты.
  2. Если pinned_count > etalon_count, дисбаланс невозможно устранить, так как rebalancer не может переместить закрепленные сегменты из этого набора реплик. В таком случае эталонное число обновляется и приравнивается к числу закрепленных сегментов. Наборы реплик с pinned_count > etalon_count не обрабатываются балансировщиком, а число закрепленных сегментов вычитается из общего числа сегментов. Балансировщик пытается переместить как можно больше сегментов из таких наборов реплик.
  3. Эта процедура повторяется с шага 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 do        can_reach_balance = true        cluster_calculate_perfect_balance(cluster, bucket_count);        foreach replicaset in cluster do                if replicaset.perfect_bucket_count <                   replicaset.pinned_bucket_count then                        can_reach_balance = false                        bucket_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 – количество наборов реплик. На каждом шаге алгоритм либо завершает вычисление, либо игнорирует хотя бы один новый набор реплик, перегруженный закрепленными сегментами, и обновляет эталонное число сегментов в других наборах реплик.

Ссылка в сегменте

Ссылка в сегменте – это счетчик в оперативной памяти, который похож на закрепление сегмента со следующими отличиями:

  1. Ссылка на сегмент не является постоянной. Ссылки предназначены для запрета перемещения сегмента во время выполнения запроса, но при перезапуске все запросы сбрасываются.

  2. Существует два типа ссылок на сегменты: только для чтения (RO) и для чтения-записи (RW).

    Если в сегменте есть ссылки типа RW, его нельзя перемещать. Однако, если балансировщику требуется отправка этого сегмента, он блокирует его для новых запросов на запись, ожидает завершения всех текущих запросов, а затем отправляет сегмент.

    Если в сегменте есть ссылки типа RO, его можно отправить, но нельзя удалить. Такой сегмент может даже перейти в статус мусора GARBAGE или отправки SENT, но его данные сохраняются до тех пор, пока не уйдет последний читатель.

    В одном сегменте могут быть ссылки как типа RO, так и типа RW.

  3. В ссылку на сегмент включён счетчик.

Методы 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: active    ref_rw: 1    ref_ro: 1    ro_lock: true    rw_lock: true    id: 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 (начиная с vshard 0.1.28). Параметры request_timeout и timeout необходимо передавать вместе, при этом должно соблюдаться следующее условие:

    timeout > request_timeout

    Например, если timeout = 10 и request_timeout = 2, в течение 10 секунд роутер может сделать 5 попыток (по 2 секунды каждая) отправки запроса разным репликам, пока запрос не завершится успешно.

  • Запросы на запись (vshard.router.callrw()) обычно не могут быть выполнены повторно без проверки того, что они не были применены ранее. Отсутствие такой проверки может привести к появлению дубликатов записей или непредвиденным изменениям данных.

    Например, клиент отправил запрос на сервер и ожидает ответ в течение заданного таймаута. Если сервер отправляет успешный ответ по истечении этого времени, клиент не получит этот ответ из-за таймаута и сочтет запрос неудачным. При повторном выполнении этого запроса без дополнительной проверки операция может быть применена дважды.

    Запрос на запись может выполняться повторно без проверки в двух случаях:

    • Запрос идемпотентен.

    • Достоверно известно, что предыдущий запрос завершился ошибкой до выполнения операций записи. Например, сервер вернул ошибку ER_READONLY. В этом случае известно, что запрос не мог быть выполнен из-за того, что сервер находится в режиме только для чтения.

Примеры дедупликации

Чтобы обеспечить идемпотентность запросов на запись (INSERT, UPDATE, UPSERT и автоинкремент), следует реализовать проверку того, что запрос применяется впервые.

Например, при добавлении нового кортежа в спейс можно использовать уникальный идентификатор вставки для проверки запроса. В примере ниже в рамках одной транзакции:

  1. Проверяется, существует ли кортеж с идентификатором key в спейсе bands.
  2. Если кортежа с таким идентификатором в спейсе нет, кортеж вставляется.
box.begin()if box.space.bands:get{key} == nil then    box.space.bands:insert{key, value}endbox.commit()

Для запросов на обновление и upsert можно создать спейс дедупликации, в котором будут сохраняться идентификаторы запросов. Спейс дедупликации – это пользовательский спейс, содержащий список уникальных идентификаторов. Каждый идентификатор соответствует одному примененному запросу. Этот спейс может иметь любое имя, в примере он называется deduplication.

В примере ниже в рамках одной транзакции:

  1. Проверяется, существует ли идентификатор запроса deduplication_key в спейсе deduplication.
  2. Если такого идентификатора нет, он добавляется в спейс дедупликации.
  3. Если запрос не был применен ранее, указанное поле в спейсе bands увеличивается на единицу. Такой подход гарантирует, что каждый запрос на изменение данных будет выполнен только один раз.
function update_1(deduplication_key, key)    box.begin()    if box.space.deduplication:get{deduplication_key} == nil then        box.space.deduplication:insert{deduplication_key}        box.space.bands:update(key, {{'+', 'value', 1 }})    end    box.commit()end

Обслуживание шардированного кластера

Сбой мастера

В случае отказа мастера в наборе реплик рекомендуется:

  1. Переключите одну из реплик в режим мастера. Это позволит новому мастеру обрабатывать все входящие запросы.
  2. Обновите конфигурацию всех участников кластера. Это перенаправит все запросы новому мастеру.

Сбой набора реплик

В случае сбоя всего набора реплик часть набора данных становится недоступной. При этом маршрутизатор пытается переподключиться к мастеру отказавшего набора реплик. Таким образом, после того как набор реплик снова заработает, кластер восстанавливается автоматически.

Плановая остановка мастера

Для проведения запланированной остановки мастера в наборе реплик рекомендуется:

  1. Обновите конфигурацию, чтобы использовать другой экземпляр в качестве мастера.
  2. Перезагрузите конфигурацию на всех экземплярах. После этого все запросы будут перенаправлены новому мастеру.
  3. Остановите старый мастер.

Плановая остановка набора реплик

Для проведения запланированной остановки набора реплик рекомендуется:

  1. Перенесите все сегменты на другие хранилища кластера. Это можно сделать, назначив набору реплик нулевой вес, что инициирует перенос его сегментов на оставшиеся узлы кластера.
  2. Обновите конфигурацию всех узлов.
  3. Остановите набор реплик.

Файберы

Поиск сегментов, восстановление сегментов и балансировка сегментов выполняются автоматически и не требуют ручного вмешательства.

С технической точки зрения есть несколько файберов, которые отвечают за различные типы действий:

  • Файбер обнаружения на маршрутизаторе выполняет поиск сегментов в фоновом режиме.
  • Файбер отказоустойчивости на маршрутизаторе поддерживает соединения с репликами.
  • Файбер сборщика мусора на каждом мастер-хранилище удаляет содержимое перемещенных сегментов.
  • Файбер восстановления сегментов на каждом мастер-хранилище восстанавливает сегменты в состояниях SENDING и RECEIVING в случае перезагрузки.
  • Файбер балансировщика на одном мастер-хранилище среди всех наборов реплик выполняет процесс балансировки. Для получения подробной информации см. разделы Процесс балансировки и Миграция сегментов.

Сборщик мусора

Файбер сборщик мусора работает в фоновом режиме на мастер-хранилищах в каждом наборе реплик. Он начинает удалять содержимое сегмента в состоянии мусора GARBAGE по частям. Когда сегмент пуст, запись о нем удаляется из системного спейса _bucket.

Восстановление сегмента

Файбер восстановления сегмента работает на мастер-хранилищах. Он помогает восстановить сегменты в статусах отправки SENDING и получения RECEIVING в случае перезагрузки.

Сегменты в статусе SENDING восстанавливаются следующим образом:

  1. Система сначала выполняет поиск сегментов в состоянии SENDING.
  2. Если такой сегмент найден, система отправляет запрос в целевой набор реплик.
  3. Если сегмент на целевом наборе реплик находится в состоянии ACTIVE, исходный сегмент удаляется с узла-источника.

Сегменты в статусе RECEIVING удаляются без дополнительных проверок.

Восстановление после отказа

Файбер отказоустойчивости работает на каждом маршрутизаторе. Если мастер набора реплик становится недоступным, файбер отказоустойчивости перенаправляет запросы на чтение к репликам. Запросы на запись отклоняются с ошибкой, пока мастер не станет доступным.