Создание шардированного кластера
Пример на GitHub: sharded_cluster_crud
В этом руководстве описывается, как запустить шардированный кластер на локальной машине и управлять им с помощью утилиты tt. В этом кластере используются следующие внешние модули:
-
vshard обеспечивает шардирование в кластере.
-
crud обеспечивает управление данными в шардированном кластере. Кластер, создаваемый в этом руководстве, включает 5 экземпляров: один роутер и 4 хранилища, которые составляют два набора реплик.

Перед началом работы:
Команда tt create используется для создания приложения на
основе предопределённого или пользовательского шаблона. Например, с
помощью встроенного шаблона vshard_cluster можно создать готовое к
запуску приложение шардированного кластера.
В этом руководстве структура приложения подготавливается вручную:
- Создайте окружение tt в текущем каталоге, выполнив команду tt init.
- В пустом каталоге
instances.enabledсозданного окруженияttсоздайте каталогsharded_cluster_crud. - В каталоге
instances.enabled/sharded_cluster_crudсоздайте следующие файлы:instances.yml–- задаёт экземпляры для запуска в текущем окружении.config.yaml–- задаёт конфигурацию кластера.storage.lua–- содержит код для хранилищ.router.lua–- содержит код для роутера.sharded_cluster_crud-scm-1.rockspec–- задаёт внешние зависимости, необходимые приложению.
В следующем разделе Разработка приложения показано, как настроить кластер и написать код для маршрутизации запросов на чтение и запись к разным хранилищам.
Откройте файл instances.yml и добавьте следующее содержимое:
storage-a-001:storage-a-002:storage-b-001:storage-b-002:router-a-001:
Этот файл задаёт экземпляры для запуска в текущем окружении.
В этом разделе описывается настройка кластера в файле config.yaml.
Добавьте раздел конфигурации credentials:
credentials:users:replicator:password: 'topsecret'roles: [ replication ]storage:password: 'secret'roles: [ sharding ]
В этом разделе создаются два пользователя с указанными паролями:
- Пользователь
replicatorс рольюreplication. - Пользователь
storageс рольюsharding.
Эти пользователи предназначены для обеспечения репликации и шардирования в кластере.
Добавьте раздел iproto.advertise:
iproto:advertise:peer:login: replicatorsharding:login: storage
В этом разделе настраиваются следующие параметры:
iproto.advertise.peer–- задаёт способ объявления текущего экземпляра другим членам кластера. В частности, этот параметр сообщает другим членам набора реплик, что для подключения к текущему экземпляру следует использовать пользователяreplicator.iproto.advertise.sharding–- задаёт способ объявления текущего экземпляра роутеру и балансировщику.
Топология кластера, определяемая в следующем разделе,
также задаёт параметр iproto.advertise.client для каждого экземпляра.
Этот параметр принимает URI, используемый для объявления экземпляра
клиентам. Например, Tarantool Cluster Manager использует эти URI
для подключения к экземплярам кластера.
Задайте общее количество бакетов в шардированном кластере с помощью параметра sharding.bucket_count:
sharding:bucket_count: 1000
Определите топологию кластера в разделе groups. Кластер включает две группы:
-
storagesвключает два набора реплик. Каждый набор реплик содержит два экземпляра. -
routersвключает один экземпляр роутера. Ниже представлена схема топологии кластера:
groups:storages:replicasets:storage-a:# ...storage-b:# ...routers:replicasets:router-a:# ...
Чтобы настроить хранилища, добавьте следующий код в раздел groups:
storages:roles: [ roles.crud-storage ]app:module: storagesharding:roles: [ storage ]replication:failover: manualreplicasets:storage-a:leader: storage-a-001instances:storage-a-001:iproto:listen:- uri: '127.0.0.1:3302'advertise:client: '127.0.0.1:3302'storage-a-002:iproto:listen:- uri: '127.0.0.1:3303'advertise:client: '127.0.0.1:3303'storage-b:leader: storage-b-001instances:storage-b-001:iproto:listen:- uri: '127.0.0.1:3304'advertise:client: '127.0.0.1:3304'storage-b-002:iproto:listen:- uri: '127.0.0.1:3305'advertise:client: '127.0.0.1:3305'
Основные параметры на уровне группы:
-
roles: этот параметр включает рольroles.crud-storage, предоставляемую модулем CRUD, для всех экземпляров хранилищ. -
app: параметрapp.moduleуказывает, что код для хранилищ должен загружаться из модуляstorage. Это описано ниже в разделе Добавление кода хранилища. -
sharding: параметр sharding.roles указывает, что все экземпляры в этой группе выступают в роли хранилищ. Балансировщик выбирается автоматически из двух master-экземпляров. -
replication: параметр replication.failover указывает, что лидер в каждом наборе реплик должен быть задан вручную. -
replicasets: в этом разделе настраиваются два набора реплик, составляющих хранилища кластера.
Чтобы настроить роутер, добавьте следующий код в раздел groups:
routers:roles: [ roles.crud-router ]app:module: routersharding:roles: [ router ]replicasets:router-a:instances:router-a-001:iproto:listen:- uri: '127.0.0.1:3301'advertise:client: '127.0.0.1:3301'
Основные параметры на уровне группы:
roles: этот параметр включает рольroles.crud-router, предоставляемую модулем CRUD, для экземпляра роутера.app: параметрapp.moduleуказывает, что код для роутера должен загружаться из модуляrouter. Это описано ниже в разделе Добавление кода роутера.sharding: параметр sharding.roles указывает, что экземпляр в этой группе выступает в роли роутера.replicasets: в этом разделе настраивается набор реплик с одним экземпляром роутера.
Полученный файл config.yaml должен выглядеть следующим образом:
credentials:users:replicator:password: 'topsecret'roles: [ replication ]storage:password: 'secret'roles: [ sharding ]
iproto:advertise:peer:login: replicatorsharding:login: storage
sharding:bucket_count: 1000
storages:roles: [ roles.crud-storage ]app:module: storagesharding:roles: [ storage ]replication:failover: manualreplicasets:storage-a:leader: storage-a-001instances:storage-a-001:iproto:listen:- uri: '127.0.0.1:3302'advertise:client: '127.0.0.1:3302'storage-a-002:iproto:listen:- uri: '127.0.0.1:3303'advertise:client: '127.0.0.1:3303'storage-b:leader: storage-b-001instances:storage-b-001:iproto:listen:- uri: '127.0.0.1:3304'advertise:client: '127.0.0.1:3304'storage-b-002:iproto:listen:- uri: '127.0.0.1:3305'advertise:client: '127.0.0.1:3305'
routers:roles: [ roles.crud-router ]app:module: routersharding:roles: [ router ]replicasets:router-a:instances:router-a-001:iproto:listen:- uri: '127.0.0.1:3301'advertise:client: '127.0.0.1:3301'
roles: [ roles.crud-router ]app:module: routersharding:roles: [ router ]replicasets:router-a:instances:router-a-001:iproto:listen:- uri: '127.0.0.1:3301'advertise:client: '127.0.0.1:3301'
roles: [ roles.crud-storage ]app:module: storagesharding:roles: [ storage ]replication:failover: manualreplicasets:storage-a:leader: storage-a-001instances:storage-a-001:iproto:listen:- uri: '127.0.0.1:3302'advertise:client: '127.0.0.1:3302'storage-a-002:iproto:listen:- uri: '127.0.0.1:3303'advertise:client: '127.0.0.1:3303'storage-b:leader: storage-b-001instances:storage-b-001:iproto:listen:- uri: '127.0.0.1:3304'advertise:client: '127.0.0.1:3304'storage-b-002:iproto:listen:- uri: '127.0.0.1:3305'advertise:client: '127.0.0.1:3305'routers:roles: [ roles.crud-router ]app:module: routersharding:roles: [ router ]replicasets:router-a:instances:router-a-001:iproto:listen:- uri: '127.0.0.1:3301'advertise:client: '127.0.0.1:3301'
Откройте файл storage.lua и определите спейс и индексы внутри box.watch() следующим образом:
box.watch('box.status', function()if box.info.ro thenreturnendbox.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)
- Функция box.schema.create_space() создает
спейс. Обратите внимание, что созданный спейс
bandsсодержит полеbucket_id. Это поле представляет ключ шардирования, используемый для распределения набора данных по разным экземплярам хранилищ.
- space_object:create_index() создает два
индекса по полям
idиbucket_id.
Откройте файл router.lua и загрузите модуль vshard следующим образом:
local vshard = require('vshard')
Откройте файл sharded_cluster_crud-scm-1.rockspec и добавьте следующее
содержимое:
package = 'sharded_cluster_crud'version = 'scm-1'source = {url = '/dev/null',}dependencies = {'vshard == 0.1.27','crud == 1.5.2'}build = {type = 'none';}
Раздел dependencies включает указанные версии модулей vshard и crud.
Чтобы установить зависимости, необходимо собрать приложение.
В терминале перейдите в каталог окружения tt. Затем выполните
команду tt build:
$ tt build sharded_cluster_crud• Running rocks makeNo existing manifest. Attempting to rebuild...• Application was successfully built
При этом модули vshard и crud, указанные в файле *.rockspec,
устанавливаются в каталог .rocks.
Чтобы запустить все экземпляры в кластере, выполните команду tt start:
$ tt start sharded_cluster_crud• Starting an instance [sharded_cluster_crud:storage-a-001]...• Starting an instance [sharded_cluster_crud:storage-a-002]...• Starting an instance [sharded_cluster_crud:storage-b-001]...• Starting an instance [sharded_cluster_crud:storage-b-002]...• Starting an instance [sharded_cluster_crud:router-a-001]...
После запуска экземпляров необходимо выполнить начальную загрузку кластера:
-
Подключитесь к экземпляру роутера с помощью
tt connect:$ tt connect sharded_cluster_crud:router-a-001• Connecting to the instance...• Connected to sharded_cluster_crud:router-a-001 -
Вызовите vshard.router.bootstrap() выполнения начальной загрузки кластера и распределения всех бакетов по наборам реплик:
sharded_cluster_crud:router-a-001> vshard.router.bootstrap()---- true
Чтобы проверить статус кластера, выполните vshard.router.info() на роутере:
sharded_cluster_crud::router-a-001> vshard.router.info()---- replicasets:storage-b:replica:network_timeout: 0.5status: availableuri: storage@127.0.0.1:3305name: storage-b-002bucket:available_rw: 500master:network_timeout: 0.5status: availableuri: storage@127.0.0.1:3304name: storage-b-001name: storage-bstorage-a:replica:network_timeout: 0.5status: availableuri: storage@127.0.0.1:3303name: storage-a-002bucket:available_rw: 500master:network_timeout: 0.5status: availableuri: storage@127.0.0.1:3302name: storage-a-001name: storage-abucket:unreachable: 0available_ro: 0unknown: 0available_rw: 1000status: 0alerts:...
Вывод включает следующие разделы:
replicasets: содержит информацию о хранилищах и их доступности.bucket: отображает общее количество бакетов для чтения-записи и только для чтения, доступных в данный момент для этого роутера.status: число от 0 до 3, показывающее наличие проблем в кластере. Значение 0 означает отсутствие проблем.alerts: может содержать описание конкретных проблем, связанных с начальной загрузкой кластера, например, проблем с подключением, событий failover или неидентифицированных бакетов.
- Чтобы вставить тестовые данные, вызовите
crud.insert_many()на роутере:
crud.insert_many('bands', {{ 1, box.NULL, 'Roxette', 1986 },{ 2, box.NULL, 'Scorpions', 1965 },{ 3, box.NULL, 'Ace of Base', 1987 },{ 4, box.NULL, 'The Beatles', 1960 },{ 5, box.NULL, 'Pink Floyd', 1965 },{ 6, box.NULL, 'The Rolling Stones', 1962 },{ 7, box.NULL, 'The Doors', 1965 },{ 8, box.NULL, 'Nirvana', 1987 },{ 9, box.NULL, 'Led Zeppelin', 1968 },{ 10, box.NULL, 'Queen', 1970 }})
Вызов этой функции распределяет данные равномерно по узлам кластера.
- Чтобы получить кортеж по указанному ID, вызовите функцию
crud.get():
sharded_cluster_crud:router-a-001> crud.get('bands', 4)---- rows:- [4, 161, 'The Beatles', 1960] metadata: [{'name': 'id','type': 'unsigned'}, {'name': 'bucket_id', 'type':'unsigned'}, {'name': 'band_name', 'type': 'string'},{'name': 'year', 'type': 'unsigned'}] - null
- Чтобы вставить новый кортеж, вызовите
crud.insert():
sharded_cluster_crud:router-a-001> crud.insert('bands', {11, box.NULL, 'The Who', 1962})---- rows:- [11, 652, 'The Who', 1962] metadata: [{'name': 'id','type': 'unsigned'}, {'name': 'bucket_id', 'type':'unsigned'}, {'name': 'band_name', 'type': 'string'},{'name': 'year', 'type': 'unsigned'}] - null
Чтобы проверить, как данные распределены по наборам реплик, выполните следующие шаги:
- Подключитесь к любому хранилищу в наборе реплик
storage-a:
$ tt connect sharded_cluster_crud:storage-a-001• Connecting to the instance...• Connected to sharded_cluster_crud:storage-a-001
Затем выберите все кортежи в спейсе bands:
sharded_cluster_crud:storage-a-001> box.space.bands:select()---- - [1, 477, 'Roxette', 1986]- [2, 401, 'Scorpions', 1965]- [4, 161, 'The Beatles', 1960]- [5, 172, 'Pink Floyd', 1965]- [6, 64, 'The Rolling Stones', 1962]- [8, 185, 'Nirvana', 1987]
- Подключитесь к любому хранилищу в наборе реплик
storage-b:
$ tt connect sharded_cluster_crud:storage-b-001• Connecting to the instance...• Connected to sharded_cluster_crud:storage-b-001
Выберите все кортежи в спейсе bands, чтобы убедиться, что он содержит другое подмножество данных:
sharded_cluster_crud:storage-b-001> box.space.bands:select()---- - [3, 804, 'Ace of Base', 1987]- [7, 693, 'The Doors', 1965]- [9, 644, 'Led Zeppelin', 1968]- [10, 569, 'Queen', 1970]- [11, 652, 'The Who', 1962]