Рекомендации по развертыванию и конфигурации кластера в нескольких ЦОД#

Введение#

Развертывание единого кластера DataGrid в нескольких центрах обработки данных (далее ЦОД) повышает гарантии сохранности данных и их доступность. Такой кластер (единый кластер DataGrid в нескольких ЦОД) называется растянутым.

Если изначально данные реплицировались между ЦОД с помощью механизма CDC и хранились на отдельных кластерах, переход на единый растянутый кластер позволит сэкономить ресурсы и упростить эксплуатацию системы (так как настроить один растянутый кластер проще, чем настраивать два обычных кластера и репликацию между ними). Но если исходная система работала в одном ЦОД, добавление второго ЦОД и миграция на растянутый кластер сделают эксплуатацию системы сложнее.

С точки зрения функциональности растянутый кластер работает идентично обычному кластеру DataGrid. Однако эффективность выполнения отдельных операций, например чтения данных, можно улучшить за счет дополнительной конфигурации кластера. Новые параметры конфигурации и примеры вариантов использования описаны в разделах ниже.

Важные ограничения по эксплуатации

Развертывать в нескольких ЦОД на данный момент можно только кластеры, которые содержат не больше 6–8 узлов. Кластеры с большим количеством узлов в растянутом режиме могут работать нестабильно (TcpDiscoverySpi сейчас находится в процессе оптимизации). Также функциональность растянутого кластера на данный момент не подходит для обслуживания критически важных систем.

Производительность low-latency систем (системы, которые требуют быстрого ответа) на растянутом кластере зависит от фактических сетевых задержек между ЦОД и от конфигурации кэшей. В режиме растянутого кластера производительность low-latency систем может значительно ухудшиться, поэтому перед использованием растянутого кластера рекомендуется отдельно изучить конфигурации и режимы работы каждой системы.

Общее описание функциональности#

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

  1. Потеря сегмента кластера в произвольном ЦОД (вследствие сетевой недоступности или отключения ЦОД) не должна приводить к потере данных. Это означает, что как минимум одна копия данных должна находиться в каждом ЦОД.

    Выполнение данного требования может быть реализовано с помощью стандартной аффинити-функции RendezvousAffinityFunction с дополнительной конфигурацией — аффинити-фильтром. Можно выбрать один из трех фильтров: MdcAffinityBackupFilter, ClusterNodeAttributeAffinityBackupFilter и фильтр с ячейками ClusterNodeAttributeColocatedBackupFilter. Они отличаются по конфигурации и по логике распределения партиций при изменении топологии. Подробнее о настройке, преимуществах и недостатках фильтров написано в разделах ниже.

  2. Компоненты DataGrid, которые требуют сетевых взаимодействий при обработке запросов пользователей, сначала должны обращаться к узлам кластера из локального ЦОД. Обращение к узлам из внешних (по отношению к узлу-обработчику пользовательского запроса) ЦОД допустимо только в том случае, когда ни один узел из локального ЦОД не подходит для выполнения запроса.

Метаинформация о расположении узлов кластера#

Идентификатор ЦОД, на котором запущен узел (сервер или толстый клиент), должен передаваться на узел с помощью системного свойства (SystemProperty) IGNITE_DATA_CENTER_ID при запуске узла. Например, как параметр Java-команды при запуске узла из командной строки: -DIGNITE_DATA_CENTER_ID=DC_0.

Идентификатором может быть произвольная строка, например - DC_0, DC_1. У разных ЦОД должны быть разные идентификаторы, а у узлов в пределах одного ЦОД — одинаковые.

Конфигурация функции распределения для обеспечения гарантии сохранности данных#

Чтобы в каждом ЦОД была полная копия данных, воспользуйтесь одним из способов ниже.

Использование бэкап-фильтра MdcAffinityBackupFilter#

Пример: кластер запущен в двух ЦОД, создается кеш с тремя резервными копиями партиций

Java#
ignite.getOrCreateCache(
            new CacheConfiguration<>()
                .setBackups(3)
                .setAffinity(
                    new RendezvousAffinityFunction()
                        .setAffinityBackupFilter(new MdcAffinityBackupFilter(2 /* количество ЦОД */, 3 /* количество резервных копий, такое же значение передается в `setBackups` */))));

Бэкап-фильтр MdcAffinityBackupFilter поддерживает только равномерное распределение копий между ЦОД, поэтому не все сочетания параметров (количество ЦОД и количество резервных копий) допустимы.

Примеры допустимых сочетаний:

  • Два ЦОД и три резервные копии партиций. Общее количество копий каждой партиции — четыре (три резервных копии партиции (backup) и одна основная (primary)). Четыре копии можно равномерно распределить по двум ЦОД (в каждый ЦОД будет назначено по две резервные копии).

  • Три ЦОД и две резервные копии партиций. Общее количество копий — три, в каждый ЦОД будет назначено по одной резервной копии.

Примеры недопустимых сочетаний:

  • Два ЦОД и две резервные копии партиций. Общее количество копий (три) невозможно равномерно распределить между двумя ЦОД.

  • Три ЦОД и одна резервная копия. Общее количество копий (две) меньше количества ЦОД.

Общее правило: полное количество копий (все резервные копии партиций +1) делится без остатка на количество ЦОД.

Использование бэкап-фильтра ClusterNodeAttributeAffinityBackupFilter#

Фильтр ClusterNodeAttributeAffinityBackupFilter при назначении резервной копии на узел проверяет, что значение хотя бы одного учитываемого при распределении партиций атрибута узла отличается от значений этого атрибута у остальных узлов, на которые уже назначены копии партиции. Атрибуты, которые анализирует фильтр, задает пользователь. Если передать фильтру атрибут ЦОД DC_ID, который всегда доступен в растянутом кластере, фильтр сможет распределять копии партиций между ЦОД.

Пример: кластер запущен в двух ЦОД, создается кеш с одной резервной копией партиции

Java#
ignite.getOrCreateCache(
    new CacheConfiguration<>()
        .setBackups(1)
        .setAffinity(
            new RendezvousAffinityFunction()
                .setAffinityBackupFilter(new ClusterNodeAttributeAffinityBackupFilter("DC_ID"))));

Из-за того, что у всех узлов в одном ЦОД одинаковое значение атрибута DC_ID, фильтр может назначить в каждый ЦОД только по одной копии партиции. Чтобы фильтр мог назначить больше одной копии в каждый ЦОД, на узлы нужно добавить дополнительные атрибуты.

В примере ниже на каждый узел добавляется дополнительный атрибут NODE_GROUP, и в каждом ЦОД запускается один узел со значением этого атрибута GROUP_A и один узел со значением GROUP_B (суммарно четыре узла в кластере). Таким образом фильтр может назначить по две копии каждой партиции в каждый ЦОД:

Пример

Java#
ignite.getOrCreateCache(
    new CacheConfiguration<>()
        .setBackups(3)
        .setAffinity(
            new RendezvousAffinityFunction()
                .setAffinityBackupFilter(new ClusterNodeAttributeAffinityBackupFilter("DC_ID", "NODE_GROUP"))));

Использование фильтра с ячейками ClusterNodeAttributeColocatedBackupFilter#

Фильтр ClusterNodeAttributeColocatedBackupFilter требует настройки одного дополнительного атрибута на каждом узле. Группа узлов с одинаковым значением такого атрибута называется ячейкой.

В фильтре с ячейками все копии одной партиции назначаются в одну ячейку, поэтому использовать параметр DC_ID в качестве атрибута для организации ячеек не получится: в этом случае все копии каждой партиции будут назначены в один ЦОД и требование «копия данных в каждом ЦОД» не будет выполнено. В растянутом кластере узлы с одинаковым значением атрибута должны присутствовать в каждом ЦОД, формируя растянутые ячейки — именно такая конфигурация позволяет добиться гарантий сохранности данных.

Пример: кластер из восьми узлов запущен в двух ЦОД, в кластере организованы две ячейки, создается кеш с тремя резервными копиями партиций

Ячейки формируются с помощью атрибута CELL. В каждом ЦОД есть два узла со значением атрибута CELL_0 и два узла со значением атрибута CELL_1:

Java#
ignite.getOrCreateCache(
            new CacheConfiguration<>()
                .setBackups(3)
                .setAffinity(
                    new RendezvousAffinityFunction()
                        .setAffinityBackupFilter(new ClusterNodeAttributeColocatedBackupFilter("CELL"))));

Сравнение фильтров#

Фильтр MdcAffinityBackupFilter#

Преимущества:

  • Более простая конфигурация — не требует дополнительной конфигурации узлов, кроме атрибута DC_ID.

  • Автоматически поддерживает количество копий в каждом ЦОД. Эта функциональность может одновременно быть и недостатком, так как может приводить к ребалансировке.

Недостатки:

  • Поддерживает только равномерное распределение — одинаковое количество копий партиции в каждом ЦОД.

  • Может приводить к запуску ребалансировки при выходе узлов (для in-memory сценария или при разрешенном baseline.autoAdjust).

Фильтр ClusterNodeAttributeAffinityBackupFilter#

Преимущества:

  • Дает точный контроль размещения партиций за счет управления атрибутами и их значениями на узлах.

  • Позволяет избежать запуска ребалансировки при выходе узла и смене базовой топологии (baseline topology).

  • Для кешей с одной резервной копией можно использовать атрибут DC_ID. В этом случае дополнительная конфигурация атрибутов на узлах не нужна.

Недостатки:

  • При выходе узлов из топологии кластера количество партиций в отдельном ЦОД может снижаться до одной, если значения атрибутов на оставшихся узлах не позволяют назначить в этот ЦОД дополнительные резервные копии.

  • Требует дополнительной конфигурации узлов (для схемы «2+ резервные копии в каждом ЦОД» требуются дополнительные атрибуты).

Фильтр с ячейками ClusterNodeAttributeColocatedBackupFilter#

Преимущества:

  • Обладает всеми преимуществами ClusterNodeAttributeAffinityBackupFilter.

  • Более высокая степень защиты от потери данных при выходе узлов: потеря происходит только при выходе всех узлов одной ячейки. При использовании ClusterNodeAttributeAffinityBackupFilter выход того же количества узлов с высокой вероятностью приведет к потере некоторых партиций.

Недостатки:

  • Требует обязательной конфигурации дополнительного атрибута для организации ячеек, переиспользовать атрибут DC_ID не получится.

  • Как и ClusterNodeAttributeAffinityBackupFilter, при выходе узлов не позволяет назначить дополнительные копии партиций, если в ЦОД нет узлов с подходящими значениями атрибутов (даже если остались другие активные узлы).

Конфигурация TopologyValidator для обработки сценария разрыва связи между ЦОД#

При разрыве связи между ЦОД растянутый кластер разъединяется на набор изолированных сегментов (независимых кластеров). Узлы DataGrid в каждом сегменте выполняют процедуру валидации топологии локального сегмента. Данная процедура нужна для определения, могут ли узлы сегмента выполнять операции модификации данных или должны перейти в режим только для чтения (read-only). Процедуру валидации выполняет компонент MdcTopologyValidator, который конфигурируется для каждого кеша. Подробнее о настройке MdcTopologyValidator написано ниже.

Если сегмент перешел в режим read-only, он останется в нем даже при восстановлении связи между ЦОД, так как объединение сегментов в общий кластер не предусмотрено. Чтобы вернуть узлы из read-only-сегмента в полнофункциональный режим работы, полностью остановите сегмент и запустите его заново. В этом случае узлы присоединятся к полнофункциональному кластеру и сами станут полнофункциональными.

Пример: MdcTopologyValidator, который сконфигурирован с основным ЦОД (вариант для четного количества ЦОД)

Java#
MdcTopologyValidator topologyValidator = new MdcTopologyValidator();
topologyValidator.setMainDatacenter("DC_0");
topologyValidator.setDatacenters(new HashSet<>(Arrays.asList("DC_0", "DC_1")));
 
ignite.getOrCreateCache(
    new CacheConfiguration<>()
        .setBackups(1)
        .setTopologyValidator(topologyValidator));

Пример: MdcTopologyValidator, который сконфигурирован с набором из трех ЦОД (вариант для нечетного количества ЦОД, нет необходимости указывать mainDatacenter)

Java#
MdcTopologyValidator topologyValidator = new MdcTopologyValidator();
topologyValidator.setDatacenters(new HashSet<>(Arrays.asList("DC_0", "DC_1", "DC_2" /* Значения атрибута идентификатора ЦОД `DC_ID` для каждого ЦОД. */)));
 
ignite.getOrCreateCache(
    new CacheConfiguration<>()
        .setBackups(1)
        .setTopologyValidator(topologyValidator));

Конфигурация:

  1. Свойство mainDatacenter — основной ЦОД. Это один из ЦОД, которые перечислены в списке datacenters (подробнее о нем написано ниже). Свойство mainDatacenter нужно указывать, если количество ЦОД четное (если оно нечетное, значение свойства можно оставить пустым).

  2. Свойство datacenters — список всех идентификаторов ЦОД через запятую. Список должен включать в себя все идентификаторы, с которыми запущены узлы кластера.

Алгоритм работы:

  1. Для четного количества ЦОД настраивается свойство mainDatacenter. При разрыве сети между ЦОД все узлы, для которых недоступны узлы из mainDatacenter, переходят в режим только для чтения (read-only). Из-за этого при физической потере mainDatacenter (если все узлы в нем остановились) все остальные узлы перейдут в режим только для чтения. Для восстановления их функций требуется вернуть в топологию узлы из mainDatacenter или (если основной ЦОД потерян на длительное время) изменить mainDatacenter в конфигурации и перезапустить кластер.

    В будущих релизах DataGrid планируется доработка, с которой станет возможной динамическая реконфигурация без необходимости перезапуска кластера.

  2. Для нечетного количества ЦОД свойство mainDatacenter следует оставить пустым. В этом случае узлы, которые переходят в режим только для чтения, определяются на основе большинства (majority). В режим read-only переходят те узлы, ЦОД которых остался в меньшинстве при разрыве связи между ЦОД. В сценарии разрыва связи одновременно между всеми ЦОД ни один сегмент кластера не сможет продолжить работу: они все перейдут в режим read-only.

Рекомендации по сетевым тайм-аутам и запуску кластера#

ConnectionRecoveryTimeout#

Свойство ConnectionRecoveryTimeout конфигурируется в TcpDiscoverySpi. Установите у свойства ConnectionRecoveryTimeout значение 0, чтобы отключить функции восстановления сетевого соединения. Это необходимо, так как в настоящее время функция не учитывает специфику растянутого кластера и может приводить к последовательной сегментации всех узлов кластера при разрыве связи между ЦОД.

Java-конфигурация свойства connectionRecoveryTimeout#
TcpDiscoverySpi tcpDisco = new TcpDiscoverySpi();
tcpDisco.setConnectionRecoveryTimeout(0);
 
IgniteConfiguration cfg = new IgniteConfiguration().setDiscoverySpi(tcpDisco);

// Запуск кластера с помощью `cfg`.

Процедура запуска кластера#

Рекомендуется запускать узлы кластера последовательно в каждом ЦОД: сначала запускаются все узлы в одном ЦОД, затем в следующем и так далее. Для кольцевой топологии такая процедура минимизирует количество соединений между узлами из разных ЦОД, таким образом улучшая стабильность работы кластера.

Инструменты для получения информации о топологии#

CLI (интерфейс командной строки)#

Возможности команды просмотра текущей топологии кластера:

  • получение информации о настроенных ЦОД и их идентификаторах;

  • распределение узлов по ЦОД;

  • поддержка переключения обозначения узлов (согласованный идентификатор узла (consistent ID)/имя экземпляра (instance name), идентификатор узла (node ID));

  • просмотр карты топологии и переходов между ЦОД;

  • предупреждение пользователя в случаях, если:

    • количество соединений между ЦОД неоптимально;

    • присутствуют толстые клиенты, которые подключены не к своему ЦОД.

  • поддерживается только TcpDiscoverySpi (ZookeeperSPI не поддерживается).

Пример вызова команды

[EXPERIMENTAL]
Prints detailed multi-DC cluster topology and ring diagnostics:

control.(sh|bat) --data-center print_topology [--node-format CONSISTENT_ID|ID|INSTANCE_NAME]

Parameters:
  --node-format CONSISTENT_ID|ID|INSTANCE_NAME  - Node format to print.
    CONSISTENT_ID                               - Consistent ID (default).
    ID                                          - ID.
    INSTANCE_NAME                               - Instance name.
Пример вывода оптимальной топологии
DC IDs: DC0, DC1

Nodes distribution by data centers:
Data center    Server nodes    Client nodes
DC0            2               0           
DC1            2               0           

--- Ring summary ---
DC transitions: 2 (minimum for 2 DCs)

Transitions map:
[DC0] srv0-dc0 ==>
[DC1] srv0-dc1 ==>
[DC0] srv1-dc0 ->
back to DC0

No client nodes are connected to a data center other than their own.
Пример вывода неоптимальной топологии
WARNING: DC transitions are not minimal! Actual transitions: 4. Minimum: 2.  

WARNING: The following client nodes are connected to a data center other than their own: cli0-dc1.

DC IDs: DC0, DC1

Nodes distribution by data centers:
Data center    Server nodes    Client nodes
DC0            3               1
DC1            3               1

--- Ring summary ---
DC transitions: 4

Transitions map:
[DC0] srv0-dc0 ==>
[DC1] srv1-dc1 ==>
[DC0] srv1-dc0 ->
      srv2-dc0 ==>
[DC1] srv0-dc1 ->
      srv2-dc1 ==>
back to DC0

Client nodes connected to a data center other than their own:
Client node    Client node DC    DC, node connected to
cli0-dc1       DC1               DC0  

Системное представление#

dataCenterId — поле в системном представлении SYS.NODES, позволяет в табличном виде вывести информацию о настроенных в топологии ЦОД.

Лог-файл#

Сообщения, которые выводятся в лог-файл:

  • При запуске узла выводится идентификатор ЦОД:

    >>> Data center ID: DC1
    
  • При появлении изменений в топологии кольца выводится список изменений:

    Ring topology changed. Transitions map:
    [DC0] srv-0-dc-0 ==>
    [DC1] srv-0-dc-1 ->
          srv-1-dc-1 ==>
    back to DC0
    
  • Периодически выводятся метрики (распределение узлов, количество соединений между ЦОД и так далее):

    Data centers metrics. Cross connections: 4 (not minimal).
        ^-- DC0 [servers=2, clients=0]
        ^-- DC1 [servers=2, clients=0]
    

Метрики#

Метрика io.discovery.ClientRouterNodeId — идентификатор узла, к которому подключен толстый клиент (TCP Discovery SPI, пересылка сообщений из кольца). Метрика может использоваться, например, для выяснения, подключен ли толстый клиент к своему ЦОД.