Механизм Change Data Capture (CDC) и межкластерная репликация#
Предусловия#
Установленный продукт DataGrid. Подробнее об установке продукта написано в документе «Руководство по установке».
Введение#
Change Data Capture (CDC) — это сценарий, предназначенный для асинхронной передачи измененных данных с целью их дальнейшей обработки.
Примеры сценариев для использования CDC:
потоковая передача изменений в Хранилище;
обновление поисковых индексов;
подсчет статистики (потоковые запросы);
аудит логов;
асинхронное взаимодействие со внешней системой (модерирование, запуск бизнес-процессов и прочее).
На механизме CDC основан процесс межкластерной репликации в DataGrid: в этом процессе CDC обрабатывает сегменты WAL-журнала, а затем доставляет локальные изменения к потребителям. Информация о межкластерной репликации в DataGrid находится в разделе «Межкластерная репликация в DataGrid» ниже.
DataGrid реализует CDC с помощью приложения ignite-cdc.sh и Java API.
CDC запускается на всех серверных узлах кластера DataGrid в отдельной JVM.
Алгоритм работы механизма CDC#
Механизм CDC работает как для persistence-кластеров, так и для in-memory-кластеров. Ниже представлена схема и алгоритм работы механизма CDC.

Для persistence-кластеров используется следующий алгоритм работы механизма:
При включенном механизме CDC серверный узел DataGrid создает в специальной директории
db/cdc/{consistent_id}(значение по умолчанию) жесткую ссылку на каждый сегмент архива WAL-журнала.Затем приложение
ignite-cdc.shзапускается на отдельной JVM и обрабатывает сегменты WAL-журнала, которые были перенесены в архив совсем недавно.После обработки сегмента приложением он удаляется. Само дисковое пространство освобождается при удалении обеих ссылок (в архиве WAL-журнала и в каталоге CDC).
CdcConsumer сохраняет состояние обработки в виде указателей на последние обработанные события в файлах cdc-*-state.bin. CdcConsumer может запрашивать сохранение (обновление) состояния у приложения ignite-cdc.sh. При запуске ignite-cdc.sh обработка событий будет продолжаться, начиная с последнего сохраненного состояния.
Внимание
Алгоритм работы механизма для in-memory-кластеров практически идентичен алгоритму работы CDC для persistence-кластеров. Отличие состоит в том, что в режиме in-memory узел записывает в WAL-журнал данные только тех областей данных (data region), для которых включен CDC.
Конфигурация#
Параметры конфигурации узла DataGrid#
Имя |
Описание |
Значение по умолчанию |
|---|---|---|
|
Свойство включает CDC на серверном узле для региона данных. Позволяет явно настроить регионы данных, для которых CDC включен |
|
|
Путь к каталогу CDC |
|
|
Тайм-аут для принудительной отправки сегмента WAL-журнала в архив (даже если обработка сегмента не завершена). Внимание Необходимо обязательно установить тайм-аут, иначе репликация будет происходить только при архивировании WAL-сегмента по его заполненности, что непредсказуемо по длительности. Расхождения между кластерами в таком случае могут достигать нескольких минут или десятков минут (зависит от нагрузки). Если нагрузка низкая, ротация может происходить еще реже |
|
Имя |
Описание |
Значение по умолчанию |
|---|---|---|
|
Свойство включает CDC на серверном узле для региона данных. Позволяет явно настроить регионы данных, для которых CDC включен |
false |
|
Путь к каталогу CDC |
«db/wal/cdc» |
|
тайм-аут для принудительной отправки сегмента WAL-журнала в архив (даже если обработка сегмента не завершена). |
-1 (disabled) |
Параметры конфигурации приложения CDC#
CDC конфигурируется так же, как и узел DataGrid — через XML-файл Spring:
ignite-cdc.shтребует указанияIgniteConfigurationиCdcConfigurationв одном конфигурационном файле;IgniteConfigurationиспользуется для установки общих параметров, например, пути к каталогу CDC,consistentIdузла и других;CdcConfigurationсодержит непосредственно настройки приложенияignite-cdc.sh.
Имя |
Описание |
Значение по умолчанию |
|---|---|---|
|
тайм-аут для ожидания блокировки. При запуске CDC устанавливает блокировку на каталог во избежание параллельной обработки каталога CDC-журнала другим работающим приложением |
1000 мс |
|
Время между последовательными проверками, в течение которого приложение находится в спящем режиме, если новые файлы отсутствуют |
1000 мс |
|
Свойство, определяющее, должны ли ключ и значение измененных записей быть предоставлены в бинарном формате |
true |
|
Реализация |
null |
|
Массив SPI «экспортеров» для передачи метрик приложения CDC. Подробнее в разделе Конфигурация экспортеров и включение метрик в New Metrics System |
null |
Свойства кластера#
Свойства кластера, перечисленные в таблице ниже, позволяют настраивать CDC во время работы кластера (подробнее о работе со свойствами кластера см. на сайте Apache Ignite).
Имя |
Описание |
Значение по умолчанию |
|---|---|---|
|
Принудительное отключение CDC во избежание переполнения диска. Полезно, если приложение CDC не работает в течение длительного времени. Внимание! Отключение CDC приведет к потере данных об изменениях |
false |
CDC API#
org.apache.ignite.cdc.CdcEvent
Полное описание интерфейса можно найти в официальной документации CdcEvent Interface.
Имя |
Описание |
|---|---|
|
Ключ измененной записи |
|
Значение измененной записи. Метод вернет |
|
ID кеша, на котором происходит изменение. Значение метода равно значению параметра |
|
Партиция DataGrid, в которой располагается измененная запись |
|
Флаг основного (primary) узла. Возвращает |
|
|
org.apache.ignite.cdc.CdcConsumer
Полное описание интерфейса можно найти в официальной документации CdcConsumer Interface.
org.apache.ignite.cdc.CdcConsumer — обработчик (потребитель) событий-изменений. CdcConsumer может быть реализован пользователем самостоятельно. В дистрибутиве поставляется три реализации CDC (подробнее в разделе «Межкластерная репликация в DataGrid»).
Имя |
Описание |
|---|---|
|
Вызывается один раз при запуске приложения CDC. Реестр |
|
Основной метод обработки изменений. Когда данный метод возвращает |
|
Вызывается один раз при остановке приложения CDC |
Метрики#
ignite-cdc.sh доступен тот же набор SPI для экспорта метрик, что и DataGrid. Следующие метрики предоставляются приложением (дополнительные метрики могут быть предоставлены потребителем):
Имя |
Описание |
|---|---|
|
Индекс текущего обрабатываемого сегмента WAL |
|
Индекс сегмента WAL, содержащего последнее зафиксированное состояние |
|
Зафиксированное смещение в байтах внутри сегмента WAL |
|
Временная метка (в миллисекундах), указывающая на начало обработки последнего сегмента |
|
Обрабатываемый приложением каталог с |
|
Обрабатываемый приложением каталог |
|
Каталог CDC, из которого приложение считывает данные |
|
Время обработки WAL-сегмента в миллисекундах |
Логирование#
ignite-cdc.sh использует ту же конфигурацию логирования, что и узел DataGrid. Единственное отличие — события записываются в файл ignite-cdc.log.
Жизненный цикл ignite-cdc.sh#
Внимание
В приложении ignite-cdc.sh реализован подход fail-fast. В случае любой ошибки приложение прекращает работу. Процедура перезапуска должна быть настроена средствами операционной системы.
После запуска приложение ignite-cdc.sh:
Находит необходимые общие каталоги. Используются значения из предоставленной
IgniteConfiguration.Блокирует CDC каталог.
Загружает сохраненное состояние («отметку» последнего обработанного события).
Запускает
CdcConsumer.Запускает бесконечный цикл ожидания и обработки новых доступных WAL-сегментов.
Останавливает потребитель данных в случае ошибки или получения сигнала об остановке.
Обработка пропущенных сегментов#
Механизм CDC может быть отключен вручную или путем настройки максимального размера каталога. В этом случае создание hardlink не происходит.
Внимание
Все изменения в пропущенных сегментах будут потеряны.
Если механизм CDC включен и был обнаружен пропуск в сегментах, например:
0000000000000002.wal;0000000000000010.wal;0000000000000011.wal;
то ignite-cdc.sh выдаст ошибку следующего вида:
Found missed segments. Some events are missed. Exiting! [lastSegment=2, nextSegment=10]
Примечание
Перед перезапуском CDC убедитесь, что данные, хранящиеся на разных кластерах, согласованы между собой. Если данные не согласованы, выполните восстановление согласованности с помощью команды повторной отправки, создания снепшота или другим пользовательским способом.
Чтобы исправить эту ошибку, выполните следующую команду с помощью утилиты control:
# Удаление из кластера CDC-ссылок до последнего пропуска в них.
control.(sh|bat) --cdc delete_lost_segment_links
# Удаление на заданном узле CDC-ссылок до последнего пропуска в них.
control.(sh|bat) --cdc delete_lost_segment_links --node-id node_id
Команда удалит все ссылки на сегменты перед последним пропуском.
Например, CDC был выключен несколько раз:
0000000000000002.wal;0000000000000003.wal;0000000000000008.wal;0000000000000010.wal;0000000000000011.wal.
Необходимо выполнить команду удаления CDC-ссылок, при этом будут удалены следующие ссылки:
0000000000000002.wal;0000000000000003.wal;0000000000000008.wal.
После запуска приложение ignite-cdc.sh начнет работу с сегмента 00000000000000000010.wal.
Принудительная повторная отправка всех данных кеша в CDC#
Если изменение данных в кеше произошло при отключенном механизме CDC, то эти изменения не будут зафиксированы обработчиком CDC, так как не будет создана CDC-ссылка на сегмент с этими данными. В этом случае необходимо выполнить повторную отправку данных из существующих кешей. Это важно, например, перед запуском репликации, чтобы восстановить согласованность данных в кешах между кластерами.
Примечание
Повторная отправка будет отменена если кластер не был отребалансирован или изменилась топология (добавлен или удален узел, изменена базовая топология).
Для повторной отправки данных кеша выполните следующую команду с помощью утилиты control:
control.(sh|bat) --cdc resend --caches cache1,...,cacheN
Команда итерируется по кешам и записывает основные копии записей данных в WAL-журнал для последующего захвата приложением CDC (ignite-cdc.sh).
Примечание
Отсутствуют гарантии уведомления обработчика CDC (реализации CdcConsumer) при параллельных обновлениях кеша — порядок событий, дублирование записей и прочее. Используйте интерфейс CdcEvent#version для определения версии записей (полное описание интерфейса можно найти в официальной документации CdcEvent Interface).
Межкластерная репликация в DataGrid#
Механизм межкластерной репликации основан на механизме CDC, описанном выше.
Механизм репликации асинхронный. Обеспечивается модулем расширения Change Data Capture Extension (cdc-ext).
В зависимости от архитектуры сети и требований безопасности доступны три схемы межкластерной репликации:
DataGrid to DataGrid через тонкий клиент Java (Java thin client).
DataGrid to DataGrid через клиентский узел.
DataGrid с использованием Platform V Corax.
Примечание
Все три схемы CDC поддерживают репликацию изменений TypeMapping и BinaryType.
Для разрешения возможных конфликтов изменения данных при использовании Active-Active режима репликации необходимо настроить модуль ConflictResolver для каждой схемы межкластерной репликации. Подробное описание этого модуля содержится в следующем разделе.
Репликация DataGrid to DataGrid через тонкий клиент Java (Java thin client)#
IgniteToIgniteCdcStreamer запускает тонкий клиент Java, подключающийся к удаленному кластеру-получателю. При наличии подключения к удаленному кластеру тонкий клиент Java будет реплицировать в него изменения данных, захваченные на кластере-источнике:

Внимание
Экземпляры приложения ignite-cdc.sh должны быть настроены и запущены на каждом серверном узле кластера-источника для передачи всех изменений данных.
Данная схема репликации требует возможности организации прямого соединения между двумя кластерами DataGrid.
Примечание
На текущий момент данная схема позволяет обеспечить репликацию только между двумя кластерами. Для соединения более двух кластеров можно воспользоваться схемой из раздела «Репликация DataGrid с использованием Platform V Corax».
Описание конфигурации#
Имя |
Описание |
Значение по умолчанию |
|---|---|---|
|
Набор (set) имен реплицируемых кешей |
null |
|
Конфигурация клиентского узла, реплицирующего изменения на кластер-получатель |
null |
|
Свойство отвечает за то, какие партиции будут реплицироваться: либо только основные (primary) партиции, либо все, чтобы в случае выхода основной партиции из строя изменения поступили от резервных (backup) партиций |
|
|
Максимальный размер пакета изменений, которые необходимо отправить |
|
Метрики#
Имя |
Описание |
|---|---|
|
Количество событий, примененных к кластеру-получателю |
|
Отметка времени последнего примененного события |
|
Количество обработанных изменений в |
|
Количество обработанных изменений в |
Репликация DataGrid to DataGrid через клиентский узел#
IgniteToIgniteCdcStreamer запускает клиентский узел, подключающийся к удаленному кластеру-получателю. При наличии подключения к удаленному кластеру клиентский узел будет реплицировать в него изменения данных, захваченные на кластере-источнике.

Внимание
Экземпляры приложения ignite-cdc.sh должны быть настроены и запущены на каждом серверном узле кластера-источника для передачи всех изменений данных.
Данная схема репликации требует возможности организации прямого соединения между двумя кластерами DataGrid.
Примечание
На текущий момент данная схема позволяет обеспечить репликацию только между двумя кластерами. Для соединения более двух кластеров можно воспользоваться схемой из раздела «Репликация DataGrid с использованием Platform V Corax».
Описание конфигурации#
Имя |
Описание |
Значение по умолчанию |
|---|---|---|
|
Набор (set) имен реплицируемых кешей |
null |
|
Конфигурация клиентского узла, реплицирующего изменения на кластер-получатель |
null |
|
Свойство отвечает за то, какие партиции будут реплицироваться: либо только основные (primary) партиции, либо все, чтобы в случае выхода основной партиции из строя изменения поступили от резервных (backup) партиций |
|
|
Максимальный размер пакета изменений, которые необходимо отправить |
|
Метрики#
Имя |
Описание |
|---|---|
|
Количество событий, реплицированных на кластер-получатель |
|
Отметка времени последнего примененного события |
|
Количество обработанных изменений в |
|
Количество обработанных изменений в |
Репликация DataGrid с использованием Platform V Corax#
Этот способ репликации изменений между кластерами требует настройки двух приложений:
ignite-cdc.shс классомorg.apache.ignite.cdc.kafka.IgniteToKafkaCdcStreamer, который будет захватывать изменения из кластера-источника и записывать их в топик Corax.kafka-to-ignite.sh, который будет считывать изменения из топика Corax и затем записывать их в кластер-получатель.
Внимание
Экземпляры приложения ignite-cdc.sh должны быть настроены и запущены на каждом серверном узле кластера-источника для передачи всех изменений данных.
Для обеспечения гарантий последовательного чтения данных при использовании репликации через Platform V Corax в топик метаданных должна быть только одна партиция.
Схема предполагает использование Platform V Corax (далее – Corax) для передачи изменений в другой кластер. Подойдут конфигурации стратегии доставки at-least once и exactly once. Идемпотентность потребителя гарантируется модулем conflictResolver, гарантирующим что повторно отправленное изменение не будет применено дважды. Кроме того, в отличие от схемы DataGrid to DataGrid, данная схема позволяет производить репликацию данных между несколькими кластерами (3 и более), которые могут содержаться в ЦОД (в примерах — ЦОД-1 и ЦОД-2). Репликация возможна как между кластерами, содержащимися в одном ЦОД, так и между кластерами, содержащимися в разных ЦОД. Различий в схеме репликации и ее настройке нет.

IgniteToKafkaCdcStreamer отправляет захваченные на кластере-источнике изменения в топик Platform V Corax. KafkaToIgniteCdcStreamerвычитывает изменения из топик и применяет на удаленном сервере-приемнике.
Для двусторонней репликации (Active-Active) запись событий CDC и IgniteToKafkaCdcStreamer (ignite-cdc.sh) включаются на всех кластерах. Для применения изменений на каждом приемнике необходимо запустить как минимум по одному экземпляру KafkaToIgniteCdcStreamer для каждого из топиков с событиями CDC других кластеров-источников (топик кластера-приемника вычитывать на самом кластре-приемнике не нужно). Кроме того, для разрешения конфликтных изменений необходмо настроить CacheVersionConflictResolver при помощи соответствующего плагина.
Описание конфигурации#
Конфигурация IgniteToKafkaCdcStreamer
Имя |
Описание |
Значение по умолчанию |
|---|---|---|
|
Набор (set) имен реплицируемых кешей |
null |
|
Свойства поставщика (producer) Corax |
null |
|
Имя топика Corax для событий CDC |
null |
|
Количество партиций топика Corax для событий CDC |
null |
|
Топик для репликации изменений |
null |
|
Свойство отвечает за то, какие партиции будут реплицироваться: либо только основные (primary) партиции, либо все, чтобы в случае выхода основной партиции из строя изменения поступили от резервных (backup) партиций |
|
|
Максимальный размер конкурентно создаваемых записей Corax. |
|
|
тайм-аут запроса Corax, мс |
|
Параметр kafkaRequestTimeout устанавливает, сколько IgniteToKafkaCdcStreamer будет ожидать завершения запроса KafkaProducer.
Внимание
Значение параметра kafkaRequestTimeout зависит от потребностей конкретного приложения. Если значение слишком маленькое, могут наблюдаться падения IgniteToKafkaCdcStreamer во время долгой отправки сообщений в Corax.
Параметр kafkaProperties устанавливает настройки KafkaProducer. Рекомендуется использовать отдельный файл для хранения необходимых свойств конфигурации. На этот файл можно сослаться в конфигурации IgniteToKafkaCdcStreamer.
bootstrap.servers=xxx.x.x.x:9092
request.timeout.ms=10000
<bean id="cdc.streamer" class="org.apache.ignite.cdc.kafka.IgniteToKafkaCdcStreamer">
<property name="topic" value="${send_data_kafka_topic_name}"/>
<property name="metadataTopic" value="${send_metadata_kafka_topic_name}"/>
<property name="kafkaPartitions" value="${send_kafka_partitions}"/>
<property name="caches">
<list>
<value>terminator</value>
</list>
</property>
<property name="onlyPrimary" value="false"/>
<property name="kafkaProperties" ref="kafkaProperties"/>
</bean>
<util:properties id="kafkaProperties" location="file:kafka_properties_path/kafka.properties"/>
Важно
Параметр request.timeout.ms является обязательным.
Метрики#
Имя |
Описание |
|---|---|
|
Количество событий, примененных к Corax |
|
Отметка времени последнего примененного события |
|
Количество обработанных изменений в |
|
Количество обработанных изменений в |
|
Количество байтов, отправленных в Corax |
|
Количество маркеров об обработанных изменениях в |
Приложение kafka-to-ignite.sh#
Приложение kafka-to-ignite.sh должно быть запущено рядом с кластером-получателем для обеспечения высокой скорости соединения. Приложение будет считывать события CDC из топик Corax и затем применять их к кластеру-получателю.
Внимание
В приложении kafka-to-ignite.sh реализован подход fail-fast. В случае любой ошибки приложение прекращает работу. Процедура перезапуска должна быть настроена средствами операционной системы.
Количество экземпляров приложения не обязательно должно совпадать с количеством узлов сервера-получателя, а должно выбираться с точки зрения достаточности для обработки нагрузки, создаваемой кластером-источником. Каждый экземпляр приложения настраивается для обработки указанного диапазона партиций топика, чтобы распределить нагрузку. Если запускается один экземпляр kafka-to-ignite.sh, то его необходимо настроить на весь диапазон партиций топика. KafkaConsumer для каждой партиции Corax будет создан для обеспечения равномерного получения данных.
Установка#
Компоненты CDC поставляются с дистрибутивом. Для включения CDC перенесите модуль cdc-ext из $IGNITE_HOME/libs/optional в $IGNITE_HOME/libs.
Конфигурация#
Приложение kafka-to-ignite.sh конфигурируется через классы POJO или XML-файл Spring как обычный узел DataGrid. Конфигурационный файл должен содержать следующие классы (beans), которые будут загружены во время запуска:
Один из экземпляров (bean) конфигурации клиента (узла или тонкого клиента):
IgniteConfiguration: конфигурация клиентского узла;ClientConfiguration: конфигурация тонкого клиента Java;
экземпляр
java.util.Propertiesс именемkafkaProperties: конфигурацияKafkaConsumer;org.apache.ignite.cdc.kafka.KafkaToIgniteCdcStreamerConfiguration: Опции, специфичные для приложенияkafka-to-ignite.sh.
Имя |
Описание |
Значение по умолчанию |
|---|---|---|
|
Набор (set) имен реплицируемых кешей |
null |
|
Тайм-аут запроса |
3000 |
|
Начальная партиция Corax (включительно) для топика событий CDC |
-1 |
|
Конечная партиция Corax (не включая) для топика CDC событий |
-1 |
|
Тайм-аут запроса Corax (мс) |
3000 |
|
Максимальный размер конкурентно создаваемых записей Corax. |
1024 |
|
Группа для |
|
|
Топик для репликации изменений |
null |
|
Количество потоков для выполнения консьюмерами. |
16 |
|
Имя топика Corax для событий CDC |
null |
Параметр kafkaRequestTimeout задает тайм-ауты для методов KafkaConsumer, кроме KafkaConsumer#poll.
Внимание
Перед установкой значения kafkaRequestTimeout нужно провести тестирование на конкретном приложении.
Параметр kafkaConsumerPollTimeout задает тайм-аут для метода KafkaConsumer#poll.
Внимание
Перед установкой значения kafkaConsumerPollTimeout нужно провести тестирование на конкретном приложении. Слишком большие значения параметра могут сильно повлиять на производительность приложения.
Партиции топиков Corax равномерно распределены по всем потокам обработки входящих сообщений (параметр threadCount). Каждый поток может обрабатывать только одну партицию за раз. Другие партиции не будут обработаны, пока не завершится процесс приема сообщений с текущей.
Конфигурация KafkaConsumer указывается с помощью bean kafkaProperties. Параметры Corax удобно хранить отдельно и задавать в конфигурации CDC-клиента следующим образом:
bootstrap.servers=xxx.x.x.x:xxxx
request.timeout.ms=10000
group.id=kafka-to-ignite-dc1
auto.offset.reset=earliest
enable.auto.commit=false
<util:properties id="kafkaProperties" location="file:kafka_properties_path/kafka.properties"/>
Важно
Параметр request.timeout.ms является обязательным.
Метрики#
Имя |
Описание |
|---|---|
|
Количество событий, полученных из Corax |
|
Отметка времени последнего события, полученного из Corax |
|
Количество событий, отправленных в кластер-приемник |
|
Отметка времени последнего события, примененного в кластере-приемнике |
|
Количество маркеров об обработанных изменениях в |
Логирование#
kafka-to-ignite.sh использует ту же конфигурацию логирования, что и узел DataGrid. Единственное отличие — события записываются в файл kafka-ignite-streamer.log.
Разрешение конфликтов репликации при помощи CacheVersionConflictResolver#
Отказоустойчивость#
В промышленной среде для обеспечения отказоустойчивости настройте CdcStreamer с параметром onlyPrimary=false, чтобы кластер-источник отправлял одно и то же изменение несколько раз, а именно (CacheConfiguration#backups + 1) раз.
Примечание
Реализация плагина из дистрибутива DataGrid не поддерживает слияние имеющейся версии и версии-кандидата, а только выбирает одну из них на основании алгоритма, описанного далее.
Реализация conflict resolver по умолчанию будет использоваться, если пользовательский conflict resolver не задан.
Плагин включен в дистрибутив DataGrid. Библиотеки расположены в каталоге libs/optional/ignite-cdc-ext. Приложение kafka-to-ignite.sh расположено в каталоге bin.
Конфигурация плагина CacheVersionConflictResolver#
Имя |
Описание |
Значение по умолчанию |
|---|---|---|
|
Локальный идентификатор кластера. Может принимать значения от 1 до 31 |
null |
|
Набор (set) кешей для обработки плагином |
null |
|
Опциональное поле для разрешения конфликта. |
null |
|
Пользовательское средство разрешения конфликта. Опционально. |
null |
Механизм работы модуля#
Реплицированные данные содержат дополнительную информацию, в частности, версии записей (CacheEntryVersion) из кластера-источника. Алгоритм разрешения конфликтов по умолчанию основан на сравнении версий записи и поля conflictResolveField.
Разрешение конфликтов на основе версии записи#
Этот подход обеспечивает окончательную гарантию согласованности только в том случае, если каждая запись обновляется только из одного кластера.
Алгоритм разрешения конфликтов:
Изменение из «локального» кластера (
clusterIdсоответствует текущему кластеру) всегда применяется безусловно. Любые реплицированные данные могут быть перезаписаны локально.Если имеющаяся запись и запись-кандидат из одного и того же кластера (
clusterIdсовпадает), то будет выбрано значение с бо́льшим значениемCacheEntryVersion.При невозможности разрешить конфликт изменение не принимается, что приводит к потере консистентности между кластерами. Ошибка разрешения конфликта выводится в log-файл.
Внимание
При таком подходе не происходит репликация обновлений и удалений записей из кластера-получателя на кластер-источник.
Для обеспечения максимальной производительности и во избежание ситуации, описанной в пункте 3, рекомендуется избегать пересечений по ключам при проектировании, а при наличии пересечений обоснованно выбирать поля для сравнения изменений.
Разрешение конфликтов на основе значения записи#
Этот подход обеспечивает конечную гарантию согласованности, даже если запись была обновлена из разных кластеров.
Примечание
Поле разрешения конфликта, указанное в conflictResolveField, должно содержать монотонно возрастающее значение, предоставленное пользователем. Например, идентификатор запроса или метку времени.
Внимание
При таком подходе не происходит репликация удалений записей из кластера-получателя на кластер-источник, поскольку удаления не могут быть версионированы по полю.
Алгоритм разрешения конфликтов:
Изменение из «локального» кластера (
clusterIdсоответствует текущему кластеру) всегда применяется безусловно. Любые реплицированные данные могут быть перезаписаны локально.Если имеющаяся запись и запись-кандидат из одного и того же кластера (
clusterIdсовпадает), то будет выбрано значение с бо́льшим значениемCacheEntryVersion.Если конфликт еще не разрешен и задано поле
conflictResolveField, то оно сравнивается у текущей записи и у записи-кандидата — выбирается запись с бо́льшим значениемconflictResolveField.При невозможности разрешить конфликт изменение не принимается, что приводит к потере консистентности между кластерами. Ошибка разрешения конфликта выводится в log-файл.
Пользовательские правила для разрешения конфликтов#
Пользователь может задать собственные правила для разрешения конфликтов, исходя из характера данных и операций. Это может быть полезно, когда стандартные стратегии для разрешения конфликтов неприменимы.
Выбор правильной стратегии разрешения конфликтов зависит от конкретного случая использования и требует хорошего понимания используемых данных и их использования. При выборе стратегии разрешения конфликтов следует учитывать характер транзакций, скорость изменения данных и последствия потенциальной потери или перезаписи данных.
Пользовательское средство разрешения конфликтов может быть задано с помощью conflictResolver и позволяет сравнивать или объединять конфликтные данные любым требуемым способом.
Пример настройки плагина CacheVersionConflictResolver#
Указывается в serverExampleConfig.xml.
<property name="pluginProviders">
<bean class="org.apache.ignite.cdc.conflictresolve.CacheVersionConflictResolverPluginProvider">
<property name="clusterId" value="1" />
<property name="caches">
<list>
<value>queryId</value>
</list>
</property>
<property name="conflictResolveField" value="modificationDate"/>
</bean>
</property>
Режимы репликации#
Существует два режима репликации:
Active-Active: в данном режиме прикладные приложения записывают данные в оба кластера. Для примеров настройки далее это означает, что можно реплицировать данные как из кластера ignite-1984 в кластер ignite-2029, так и наоборот;
Active-Passive: в данном режиме прикладные приложения записывают данные только в один кластер. Второй кластер в данном режиме неактивен и используется только в случае возникновения проблем (холодный резерв, standIn) или для выполнения запросов на чтение (например, для составления отчетов).
Данные режимы работают для обеих представленных схем репликации.
Гарантии механизмов CDC-репликации DataGrid#
Механизм CDC не является заменой транзакционной репликации и имеет свои особенности. Основные гарантии и ограничения, которые необходимо учитывать при использовании CDC:
Атомарность: CDC гарантирует неделимость операций
#put()и#remove(). Это значит, что захват изменений для этих операций будет происходить целиком или не будет происходить вовсе.Транзакционность: CDC не поддерживает транзакционные гарантии. Это значит, что если несколько операций изменения данных (например,
putиremove) происходят в рамках одной транзакции, CDC может зафиксировать их по отдельности, а не как одну неделимую операцию.
CDC не гарантирует синхронного обновления кластеров. Репликация изменений происходит асинхронно, что может привести к временной рассогласованности данных между кластерами.
Гарантии порядка (последовательности):
В пределах одного ключа: при работе одного экземпляра прикладного кода с одним кластером (режим Active-Passive) CDC гарантирует порядок применения событий для одного ключа. Изменения, которые произошли с одним и тем же ключом, будут обработаны в порядке их возникновения. Для двух и более экземпляров прикладного кода такая гарантия отсутствует.
В пределах нескольких ключей: гарантии последовательности отсутствуют для сообщений с разными ключами. Изменения для разных ключей могут обрабатываться в произвольном порядке.
При работе прикладного кода с двумя кластерами одновременно (режим Active-Active) гарантия последовательности применения событий механизмом CDC теряет смысл. Консистентность данных может быть нарушена из-за асинхронности репликации. В данном случае гарантии согласованности будут обеспечены модулем conflictResolver — подробнее об этом написано выше в разделе «Механизм работы модуля».
Гарантии доставки
Параметр onlyPrimary конфигурации CDC позволяет выбрать стратегию репликации для основных (primary) и резервных (backup) партиций кеша:
onlyPrimary=true— реплицируется только primary-запись для каждого ключа.onlyPrimary=false— реплицируются все существующие копии для каждого ключа.
В процессе обработки WAL-архива CDC обработчик CdcConsumer отслеживает индекс текущего WAL-сегмента. Эта информация заносится в параметр state обработчика CDC и обновляется в результате успешной обработки всего сегмента. В случае остановки и перезапуска CdcConsumer это позволит продолжить обработку очереди сегмента с актуального места.
Такие способы репликации, как DataGrid-to-DataGrid и DataGrid-to-Corax, предоставляют гарантии доставки at-least once для onlyPrimary=true и exactly once для onlyPrimary=false при полной работоспособности кластера (все узлы кластера работоспособны и находятся в топологии).
Примечание
При использовании продукта Platform V Corax (KFK) в CDC гарантии последовательности и дублирования событий при записи в кластер-приемник на этапе Corax-to-DataGrid зависят от конфигурации топика. Для последовательного чтения данных в топике метаданных должна быть только одна партиция. Дублирование событий в кластере Corax можно обеспечить настройкой соответствующей гарантии топика (At-most-once, At-least-once, Exactly-once).
Отказоустойчивость#
Для обеспечения гарантии дублирования записи при репликации можно использовать флаг onlyPrimary=false. Это повышает отказоустойчивость репликации к выходу узлов кластера из топологии.
Примечание
Кроме флага onlyPrimary=false, на гарантии дублирования записи также положительно влияют:
увеличение количества резервных копий (бэкап-фактор),
использование реплицированного кеша (режим
REPLICATED).
Чем больше копий партиции в кластере-исходнике, тем больше копий будет реплицировано в кластер-приемник.
Внимание
Дублирование увеличивает надежность репликации, но может отрицательно сказываться на ее скорости, так как репликация копий партиций требует дополнительных ресурсов. Конфигурация флага onlyPrimary сводится к выбору между надежностью и скоростью репликации.
Принудительная повторная отправка всех данных кеша в CDC#
При использовании утилиты control.sh для повторной отправки данных кеша нет гарантий относительно порядка событий, дублирования записей или уведомления обработчика CDC при параллельных обновлениях кеша. Это значит, что при массовой повторной отправке данных возможно их получение в другом порядке или получение дубликатов (в зависимости от параллелизма операций записи в кеш).
Примечание
Механизм репликации для in-memory-кластеров (online CDC) дает аналогичные гарантии по порядку и дублированию событий.
Примеры конфигурации при различных способах репликации#
DataGrid to DataGrid через клиентский узел#
Примечание
Конфигурация ниже предназначена для локального развертывания (на localhost). Настройки указаны для примера.
Конфигурация кластера-источника:
Имя файла: ignite-1984.xml.
<beans xmlns="http://www.springframework.org/schema/beans"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xmlns:util="http://www.springframework.org/schema/util"
xsi:schemaLocation="
http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd
http://www.springframework.org/schema/util http://www.springframework.org/schema/util/spring-util.xsd">
<bean class="org.apache.ignite.configuration.IgniteConfiguration">
<property name="igniteInstanceName" value="ignite-1984-server" />
<property name="consistentId" value="ignite-1984-server" />
<property name="localHost" value="xxx.x.x.x" />
<property name="discoverySpi">
<bean class="org.apache.ignite.spi.discovery.tcp.TcpDiscoverySpi">
<property name="localPort" value="47500" />
<property name="ipFinder">
<bean class="org.apache.ignite.spi.discovery.tcp.ipfinder.vm.TcpDiscoveryVmIpFinder">
<property name="addresses"
value="xxx.x.x.x:47500..47510" />
</bean>
</property>
<property name="joinTimeout" value="10000" />
</bean>
</property>
<property name="dataStorageConfiguration">
<bean class="org.apache.ignite.configuration.DataStorageConfiguration">
<property name="defaultDataRegionConfiguration">
<bean class="org.apache.ignite.configuration.DataRegionConfiguration">
<property name="persistenceEnabled" value="true" />
<!-- Свойство `cdcEnabled` (ниже) включает функциональность CDC. -->
<property name="cdcEnabled" value="true" />
</bean>
</property>
<!-- При данном значении свойства `walForceArchiveTimeout` выше каждый сегмент WAL-журнала будет доступен для обработки через 5 секунд после перенесения его в архив. -->
<property name="walForceArchiveTimeout" value="5000"/>
</bean>
</property>
<!--Ниже нужно указать конфигурацию модуля `conflictResolver`. -->
<property name="pluginProviders">
<bean class="org.apache.ignite.cdc.conflictresolve.CacheVersionConflictResolverPluginProvider">
<property name="clusterId" value="1" />
<!-- `clusterId` (выше) должен быть разным для каждого кластера. -->
<!-- Ниже нужно сконфигурировать кeши, для которых будет срабатывать `conflictResolver`. -->
<property name="caches">
<util:list>
<bean class="java.lang.String">
<constructor-arg type="String" value="terminator" />
<property name="conflictResolveField" value="modificationDate"/>
</bean>
</util:list>
</property>
</bean>
</property>
<property name="cacheConfiguration">
<list>
<bean class="org.apache.ignite.configuration.CacheConfiguration">
<property name="atomicityMode" value="ATOMIC"/>
<property name="name" value="terminator"/>
</bean>
</list>
</property>
</bean>
<!-- Ниже указывается конфигурация механизма CDC. Утилита `ignite-cdc.sh` должна запускаться на всех серверных узлах кластера DataGrid, поэтому конфигурация CDC хранится в том же файле, в котором хранится конфигурация серверного узла. -->
<bean id="cdc.cfg" class="org.apache.ignite.cdc.CdcConfiguration">
<!--Конфигурация здесь отличается от конфигурации по умолчанию только свойством `consumer` (потребитель). Здесь оно относится к классу `IgniteToIgniteCdcStreamer`. -->
<property name="consumer" ref="cdc.streamer" />
</bean>
<!-- Внутри класса передается конфигурация клиентского узла, который будет подключаться к кластеру-получателю данных. -->
<bean id="cdc.streamer" class="org.apache.ignite.cdc.IgniteToIgniteCdcStreamer">
<property name="destinationIgniteConfiguration">
<bean class="org.apache.ignite.configuration.IgniteConfiguration">
<property name="igniteInstanceName" value="ignite-2029-cdc-client" />
<property name="clientMode" value="true" />
<property name="localHost" value="xxx.x.x.x" />
<property name="discoverySpi">
<bean class="org.apache.ignite.spi.discovery.tcp.TcpDiscoverySpi">
<property name="localPort" value="47600" />
<property name="ipFinder">
<bean class="org.apache.ignite.spi.discovery.tcp.ipfinder.vm.TcpDiscoveryVmIpFinder">
<property name="addresses"
value="xxx.x.x.x:47600..47610" />
</bean>
</property>
<property name="joinTimeout" value="10000" />
</bean>
</property>
</bean>
</property>
<property name="onlyPrimary" value="false"/>
<property name="caches">
<list>
<value>terminator</value>
</list>
</property>
<property name="maxBatchSize" value="256"/>
</bean>
</beans>
Конфигурация кластера-получателя:
Примечание
Представленная ниже конфигурация актуальна для режима репликации Active-Active. При использовании режима Active-Passive на кластере-получателе org.apache.ignite.cdc.CdcConfiguration не настраивается. Также при Active-Passive на обоих кластерах не потребуется настройка conflictResolver. Для получения более подробной информации о режимах репликации перейдите к разделу «Режимы репликации» ниже.
Имя файла: ignite-2029.xml.
<beans xmlns="http://www.springframework.org/schema/beans"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xmlns:util="http://www.springframework.org/schema/util"
xsi:schemaLocation="
http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd
http://www.springframework.org/schema/util http://www.springframework.org/schema/util/spring-util.xsd">
<bean class="org.apache.ignite.configuration.IgniteConfiguration">
<property name="igniteInstanceName" value="ignite-2029" />
<property name="consistentId" value="ignite-2029" />
<property name="localHost" value="xxx.x.x.x" />
<property name="discoverySpi">
<bean class="org.apache.ignite.spi.discovery.tcp.TcpDiscoverySpi">
<property name="localPort" value="47600" />
<property name="ipFinder">
<bean class="org.apache.ignite.spi.discovery.tcp.ipfinder.vm.TcpDiscoveryVmIpFinder">
<property name="addresses"
value="xxx.x.x.x:47600..47610" />
</bean>
</property>
<property name="joinTimeout" value="10000" />
</bean>
</property>
<property name="dataStorageConfiguration">
<bean class="org.apache.ignite.configuration.DataStorageConfiguration">
<property name="defaultDataRegionConfiguration">
<bean class="org.apache.ignite.configuration.DataRegionConfiguration">
<property name="persistenceEnabled" value="true" />
<!-- Свойство `cdcEnabled` (ниже) включает функциональность CDC. -->
<property name="cdcEnabled" value="true" />
</bean>
</property>
<property name="walForceArchiveTimeout" value="5000"/>
</bean>
</property>
<property name="pluginProviders">
<bean class="org.apache.ignite.cdc.conflictresolve.CacheVersionConflictResolverPluginProvider">
<property name="clusterId" value="2" />
<property name="caches">
<util:list>
<bean class="java.lang.String">
<constructor-arg type="String" value="terminator" />
<property name="conflictResolveField" value="modificationDate"/>
</bean>
</util:list>
</property>
</bean>
</property>
<property name="cacheConfiguration">
<list>
<bean class="org.apache.ignite.configuration.CacheConfiguration">
<property name="atomicityMode" value="ATOMIC"/>
<property name="name" value="terminator"/>
</bean>
</list>
</property>
</bean>
<bean id="cdc.cfg" class="org.apache.ignite.cdc.CdcConfiguration">
<property name="consumer" ref="cdc.streamer" />
</bean>
<bean id="cdc.streamer" class="org.apache.ignite.cdc.IgniteToIgniteCdcStreamer">
<property name="destinationIgniteConfiguration">
<bean class="org.apache.ignite.configuration.IgniteConfiguration">
<property name="igniteInstanceName" value="ignite-1984-cdc-client" />
<property name="clientMode" value="true" />
<property name="localHost" value="xxx.x.x.x" />
<property name="discoverySpi">
<bean class="org.apache.ignite.spi.discovery.tcp.TcpDiscoverySpi">
<property name="localPort" value="47500" />
<property name="ipFinder">
<bean class="org.apache.ignite.spi.discovery.tcp.ipfinder.vm.TcpDiscoveryVmIpFinder">
<property name="addresses"
value="xxx.x.x.x:47500..47510" />
</bean>
</property>
<property name="joinTimeout" value="10000" />
</bean>
</property>
</bean>
</property>
<property name="onlyPrimary" value="false"/>
<property name="caches">
<list>
<value>terminator</value>
</list>
</property>
<property name="maxBatchSize" value="256"/>
</bean>
</beans>
После конфигурирования кластеров необходимо запустить утилиту ignite-cdc.sh на каждом серверном узле. Для запуска выполните следующие действия:
Выполните команду
ignite-cdc.sh ignite-1984.xml.Выполните команду
ignite-cdc.sh ignite-2029.xml.
DataGrid to DataGrid через тонкий клиент Java#
Примечание
Конфигурация ниже предназначена для локального развертывания (на localhost). Настройки указаны для примера.
Настройка аналогична настройке «Репликация DataGrid to DataGrid через клиентский узел». Отличие в том, что в экземпляре CdcConfiguration в поле consumer необходимо указать IgniteToIgniteClientStreamer.
<bean class="org.apache.ignite.cdc.CdcConfiguration">
<property name="consumer">
<bean class="org.apache.ignite.cdc.thin.IgniteToIgniteClientCdcStreamer">
<property name="destinationClientConfiguration">
<bean class="org.apache.ignite.configuration.ClientConfiguration">
<property name="addresses">
<list>
<!-- Список адресов кластера-получателя. -->
<value>xxx.x.x.x:10800</value>
</list>
</property>
</bean>
</property>
</bean>
</property>
<!-- Остальная конфигурация `CdcConfiguration`. -->
</bean>
DataGrid с использованием Platform V Corax#
Пример настройки двусторонней (Active-Active) репликации
Примечание
В примере ниже описывается конфигурация для локального развертывания (на localhost) и передачи данных между двумя ЦОД. Настройки указаны для примера.
Перед началом работы запустите Corax и Zookeeper. Затем создайте топики (число партиций и названия для примера):
bin/kafka-topics.sh --create --partitions 16 --replication-factor 1 --topic dc1_to_dc2 --bootstrap-server localhost:9092
bin/kafka-topics.sh --create --partitions 16 --replication-factor 1 --topic dc2_to_dc1 --bootstrap-server localhost:9092
bin/kafka-topics.sh --create --partitions 1 --replication-factor 1 --topic metadata_from_dc1 --bootstrap-server localhost:9092
bin/kafka-topics.sh --create --partitions 1 --replication-factor 1 --topic metadata_from_dc2 --bootstrap-server localhost:9092
Конфигурация кластера-источника для ЦОД-1
Имя файла: ignite-dc1-kafka.xml.
<beans xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xmlns:util="http://www.springframework.org/schema/util" xmlns="http://www.springframework.org/schema/beans"
xsi:schemaLocation="http://www.springframework.org/schema/beans
http://www.springframework.org/schema/beans/spring-beans.xsd
http://www.springframework.org/schema/util
https://www.springframework.org/schema/util/spring-util.xsd">
<bean class="org.apache.ignite.configuration.IgniteConfiguration">
<property name="igniteInstanceName" value="ignite-2029-server"/>
<property name="consistentId" value="ignite-2029-server"/>
<property name="cacheConfiguration">
<list>
<bean class="org.apache.ignite.configuration.CacheConfiguration">
<property name="name" value="terminator"/>
</bean>
</list>
</property>
<property name="discoverySpi">
<bean class="org.apache.ignite.spi.discovery.tcp.TcpDiscoverySpi">
<property name="localPort" value="47800"/>
<property name="ipFinder">
<bean class="org.apache.ignite.spi.discovery.tcp.ipfinder.vm.TcpDiscoveryVmIpFinder">
<property name="addresses" value="xxx.x.x.x:47800"/>
</bean>
</property>
<property name="joinTimeout" value="10000"/>
</bean>
</property>
<property name="dataStorageConfiguration">
<bean class="org.apache.ignite.configuration.DataStorageConfiguration">
<property name="defaultDataRegionConfiguration">
<bean class="org.apache.ignite.configuration.DataRegionConfiguration">
<property name="persistenceEnabled" value="true"/>
<!-- Свойство `cdcEnabled` (ниже) включает функциональность CDC. -->
<property name="cdcEnabled" value="true"/>
</bean>
</property>
<!-- При данном значении свойства `walForceArchiveTimeout` каждый сегмент WAL-журнала будет доступен для обработки через 5 секунд после перенесения его в архив. -->
<property name="walForceArchiveTimeout" value="5000"/>
</bean>
</property>
<property name="clientConnectorConfiguration" ref="clientConnectorConfiguration"/>
<property name="communicationSpi" ref="communicationSpi"/>
<property name="connectorConfiguration" ref="connectorConfiguration"/>
<!-- Ниже нужно указать конфигурацию модуля `conflictResolver`. -->
<property name="pluginProviders">
<list>
<bean class="org.apache.ignite.cdc.conflictresolve.CacheVersionConflictResolverPluginProvider">
<property name="clusterId" value="1"/>
<!-- `clusterId` (выше) должен быть разным для каждого кластера. -->
<property name="caches">
<!-- Ниже нужно сконфигурировать кеши, для которых будет срабатывать `conflictResolver`. -->
<set>
<value>terminator</value>
</set>
</property>
<property name="conflictResolveField" value="modificationDate"/>
</bean>
</list>
</property>
</bean>
<!-- Ниже указывается конфигурация сервера. -->
<bean id="clientConnectorConfiguration" class="org.apache.ignite.configuration.ClientConnectorConfiguration">
<property name="port" value="11000"/>
</bean>
<!-- Ниже указывается конфигурация сервера. -->
<bean id="connectorConfiguration" class="org.apache.ignite.configuration.ConnectorConfiguration">
<property name="port" value="11300"/>
</bean>
<!-- Ниже указывается конфигурация сервера. -->
<bean id="communicationSpi" class="org.apache.ignite.spi.communication.tcp.TcpCommunicationSpi">
<property name="localPort" value="48800"/>
</bean>
<util:properties id="kafkaProperties" location="file:/config/path/kafka.properties"/>
<bean id="cdc.cfg" class="org.apache.ignite.cdc.CdcConfiguration">
<property name="consumer">
<!-- Ниже указывается конфигурация механизма CDC. Утилита `ignite-cdc.sh` должна запускаться на всех серверных узлах кластера DataGrid, поэтому конфигурация CDC хранится в том же файле, в котором хранится конфигурация серверного узла. -->
<bean class="org.apache.ignite.cdc.kafka.IgniteToKafkaCdcStreamer">
<property name="topic" value="dc1_to_dc2"/>
<property name="metadataTopic" value="metadata_from_dc1"/>
<property name="kafkaPartitions" value="16"/>
<property name="caches">
<list>
<value>terminator</value>
</list>
</property>
<property name="maxBatchSize" value="256"/>
<property name="onlyPrimary" value="false"/>
<property name="kafkaProperties" ref="kafkaProperties"/>
</bean>
</property>
</bean>
</beans>
Конфигурация кластера-источника для ЦОД-2
Имя файла: ignite-dc2-kafka.xml.
<beans xmlns="http://www.springframework.org/schema/beans"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xmlns:util="http://www.springframework.org/schema/util"
xsi:schemaLocation="http://www.springframework.org/schema/beans
http://www.springframework.org/schema/beans/spring-beans.xsd
http://www.springframework.org/schema/util https://www.springframework.org/schema/util/spring-util.xsd">
<bean class="org.apache.ignite.configuration.IgniteConfiguration">
<property name="igniteInstanceName" value="ignite-1984-server"/>
<property name="consistentId" value="ignite-1984-server"/>
<property name="cacheConfiguration">
<list>
<bean class="org.apache.ignite.configuration.CacheConfiguration">
<property name="name" value="terminator"/>
</bean>
</list>
</property>
<property name="discoverySpi">
<bean class="org.apache.ignite.spi.discovery.tcp.TcpDiscoverySpi">
<property name="localPort" value="47850"/>
<property name="ipFinder">
<bean class="org.apache.ignite.spi.discovery.tcp.ipfinder.vm.TcpDiscoveryVmIpFinder">
<property name="addresses" value="xxx.x.x.x:47850"/>
</bean>
</property>
<property name="joinTimeout" value="10000"/>
</bean>
</property>
<property name="dataStorageConfiguration">
<bean class="org.apache.ignite.configuration.DataStorageConfiguration">
<property name="defaultDataRegionConfiguration">
<bean class="org.apache.ignite.configuration.DataRegionConfiguration">
<property name="persistenceEnabled" value="true"/>
<!-- Свойство `cdcEnabled` (ниже) включает функциональность CDC. -->
<property name="cdcEnabled" value="true"/>
</bean>
</property>
<!-- При данном значении свойства `walForceArchiveTimeout` каждый сегмент WAL-журнала будет доступен для обработки через 5 секунд после перенесения его в архив. -->
<property name="walForceArchiveTimeout" value="5000"/>
</bean>
</property>
<!-- Ниже нужно указать конфигурацию модуля `conflictResolver`. -->
<property name="pluginProviders">
<list>
<bean class="org.apache.ignite.cdc.conflictresolve.CacheVersionConflictResolverPluginProvider">
<!-- `clusterId` должен быть разным для каждого кластера. -->
<property name="clusterId" value="2"/>
<property name="caches">
<!-- Здесь нужно сконфигурировать кеши, для которых будет срабатывать `conflictResolver`. -->
<set>
<value>terminator</value>
</set>
</property>
<property name="conflictResolveField" value="modificationDate"/>
</bean>
</list>
</property>
<property name="clientConnectorConfiguration" ref="clientConnectorConfiguration"/>
<property name="communicationSpi" ref="communicationSpi"/>
<property name="connectorConfiguration" ref="connectorConfiguration"/>
</bean>
<bean id="clientConnectorConfiguration" class="org.apache.ignite.configuration.ClientConnectorConfiguration">
<property name="port" value="11050"/>
</bean>
<bean id="connectorConfiguration" class="org.apache.ignite.configuration.ConnectorConfiguration">
<property name="port" value="11350"/>
</bean>
<bean id="communicationSpi" class="org.apache.ignite.spi.communication.tcp.TcpCommunicationSpi">
<property name="localPort" value="48850"/>
</bean>
<util:properties id="kafkaProperties" location="file:/config/path/kafka.properties"/>
<bean id="cdc.cfg" class="org.apache.ignite.cdc.CdcConfiguration">
<property name="consumer">
<bean class="org.apache.ignite.cdc.kafka.IgniteToKafkaCdcStreamer">
<property name="topic" value="dc2_to_dc1"/>
<property name="metadataTopic" value="metadata_from_dc2"/>
<property name="kafkaPartitions" value="16"/>
<property name="caches">
<list>
<value>terminator</value>
</list>
</property>
<property name="maxBatchSize" value="256"/>
<property name="onlyPrimary" value="false"/>
<property name="kafkaProperties" ref="kafkaProperties"/>
</bean>
</property>
</bean>
</beans>
Пример конфигурации для подключения к Platform V Corax
Примечание
Настройки ниже являются содержимым файла kafka.properties. В данном примере его содержимое одинаково для обоих ЦОД.
bootstrap.servers=xxx.x.x.x:9092
request.timeout.ms=10000
Конфигурация kafka-to-ignite.sh в ЦОД-1, вычитывающего изменения из ЦОД-2
Примечание
Представленная ниже конфигурация актуальна для режима репликации Active-Active. При использовании режима Active-Passive на кластере-получателе org.apache.ignite.cdc.CdcConfiguration и CacheVersionConflictResolver не настраивается. Для получения более подробной информации о режимах репликации перейдите к разделу «Режимы репликации» ниже.
Имя файла: kafka2ignite_dc1.xml.
<beans xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xmlns:util="http://www.springframework.org/schema/util"
xmlns="http://www.springframework.org/schema/beans"
xsi:schemaLocation="
http://www.springframework.org/schema/beans
http://www.springframework.org/schema/beans/spring-beans.xsd
http://www.springframework.org/schema/util
http://www.springframework.org/schema/util/spring-util.xsd">
<bean id="streamer.cfg" class="org.apache.ignite.cdc.kafka.KafkaToIgniteCdcStreamerConfiguration">
<property name="topic" value="dc2_to_dc1"/>
<property name="metadataTopic" value="metadata_from_dc2"/>
<property name="kafkaPartsFrom" value="0"/>
<property name="kafkaPartsTo" value="16"/>
<property name="threadCount" value="4"/>
<property name="caches">
<list>
<value>terminator</value>
</list>
</property>
</bean>
<util:properties id="kafkaProperties" location="file:/config/path/kafka2ignite_dc1.properties"/>
<bean id="ignIgniteConfiguration" class="org.apache.ignite.configuration.IgniteConfiguration">
<property name="discoverySpi" ref="ignTcpDiscoverySpi"/>
<property name="clientMode" value="true"/>
<property name="consistentId" value="kafka-to-ignite_dc2"/>
</bean>
<!-- `TcpDiscoverySpi`. -->
<bean id="ignTcpDiscoverySpi" class="org.apache.ignite.spi.discovery.tcp.TcpDiscoverySpi">
<property name="localPort" value="47801"/>
<property name="ipFinder">
<bean class="org.apache.ignite.spi.discovery.tcp.ipfinder.vm.TcpDiscoveryVmIpFinder">
<property name="addresses">
<list>
<value>xxx.x.x.x:47800..47809</value>
</list>
</property>
</bean>
</property>
</bean>
</beans>
Конфигурация KafkaConsumer для ЦОД-1
Файл kafka2ignite_dc1.properties
bootstrap.servers=xxx.x.x.x:9092
request.timeout.ms=10000
group.id=kafka-to-ignite-dc1
auto.offset.reset=earliest
enable.auto.commit=false
Конфигурация kafka-to-ignite.sh в ЦОД-2, вычитывающего изменения из ЦОД-1
Имя файла: kafka2ignite_dc2.xml.
<beans xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xmlns:util="http://www.springframework.org/schema/util"
xmlns="http://www.springframework.org/schema/beans"
xsi:schemaLocation="
http://www.springframework.org/schema/beans
http://www.springframework.org/schema/beans/spring-beans.xsd
http://www.springframework.org/schema/util
http://www.springframework.org/schema/util/spring-util.xsd">
<bean id="streamer.cfg" class="org.apache.ignite.cdc.kafka.KafkaToIgniteCdcStreamerConfiguration">
<property name="topic" value="dc1_to_dc2"/>
<property name="metadataTopic" value="metadata_from_dc1"/>
<property name="kafkaPartsFrom" value="0"/>
<property name="kafkaPartsTo" value="16"/>
<property name="threadCount" value="4"/>
<property name="caches">
<list>
<value>terminator</value>
</list>
</property>
</bean>
<util:properties id="kafkaProperties" location="file:/config/path/kafka2ignite_dc2.properties"/>
<bean id="ignIgniteConfiguration" class="org.apache.ignite.configuration.IgniteConfiguration">
<property name="discoverySpi" ref="ignTcpDiscoverySpi"/>
<property name="clientMode" value="true"/>
<property name="consistentId" value="kafka-to-ignite_dc2"/>
</bean>
<!-- `TcpDiscoverySpi`. -->
<bean id="ignTcpDiscoverySpi" class="org.apache.ignite.spi.discovery.tcp.TcpDiscoverySpi">
<property name="localPort" value="47851"/>
<property name="ipFinder">
<bean class="org.apache.ignite.spi.discovery.tcp.ipfinder.vm.TcpDiscoveryVmIpFinder">
<property name="addresses">
<list>
<value>xxx.x.x.x:47850..47859</value>
</list>
</property>
</bean>
</property>
</bean>
</beans>
Конфигурация KafkaConsumer для ЦОД-2:
Файл kafka2ignite_dc2.properties
bootstrap.servers=xxx.x.x.x:9092
request.timeout.ms=10000
group.id=kafka-to-ignite-dc2
auto.offset.reset=earliest
enable.auto.commit=false
Конфигурация kafka-to-ignite.sh c тонким клиентом
Возможно использование тонкого клиента вместо клиентского узла.
Настройка аналогична настройке kafka-to-ignite.sh с клиентским узлом с тем отличием, что вместо экземпляра IgniteConfiguration необходимо указать ClientConfiguration.
<bean id="client.cfg" class="org.apache.ignite.configuration.ClientConfiguration">
<property name="addresses">
<list>
<!-- Список адресов кластера-получателя. -->
<value>xxx.x.x.x:10800</value>
</list>
</property>
</bean>
Остальные настройки KafkaToIgniteCdcStreamerConfiguration и kafkaProperties аналогичны.
После конфигурирования кластеров необходимо запустить утилиту ignite-cdc.sh на каждом серверном узле, а на принимающем узле необходимо запустить утилиту kafka-to-ignite.sh. В ignite-cdc.sh необходимо передавать путь к XML-конфигурации серверных узлов, а в kafka-to-ignite.sh — путь к XML-конфигурации Platform V Corax to DataGrid:
Запуск DataGrid to Platform V Corax:
ignite-cdc.sh ignite-dc1-kafka.xmlignite-cdc.sh ignite-dc2-kafka.xmlЗапуск Platform V Corax to DataGrid:
kafka-to-ignite.sh kafka2ignite_dc1.xmlkafka-to-ignite.sh kafka2ignite_dc2.xml
Возможные проблемы в работе репликации и методы их устранения#
Чтобы
CdcStreamerне завершил работу с ошибкой, до запуска репликации необходимо:Активировать кластер-получатель.
Создать реплицируемый кеш на кластере-получателе.
Во избежание возникновения сбоев на серверных узлах
cdcWalPathиwalArchivePathдолжны быть расположены в одном разделе файловой системы. В противном случае серверный узел завершит работу с ошибкой.Во избежание зацикливания сообщений при репликации в режиме Active-Active необходимо убедиться, что:
Для всех реплицируемых кешей настроен
conflictResolver.В
conflictResolverуказаны разныеclusterIdдля разных кластеров.Если conflictResolveField не указан в conflictResolver, то для уже имеющихся локально записей конфликтные изменения с удаленного кластера разрешаться не будут.
Во избежание ошибок при работе с одинаково названными SQL-таблицами необходимо настроить
value_type. По умолчанию узел-получатель генерирует свое собственное имя для связанной с ним SQL-таблицы. Чтобы реплицированные данные корректно отображались через SQL-запросы на узле-получателе, укажите одинаковое значениеvalue_typeдля обеих таблиц на исходном и принимающем узлах.Пример кода
CREATE TABLE Person (id int primary key, varchar name) with cache_name='Person', value_type='SQL_PERSON_TABLE_TYPE'При создании таблицы в DDL-команде добавьте:
имя типа данных (обязательно) — укажите уникальный численно-буквенный код;
имя кеша (опционально) — используйте любой кеш, не существующий на момент запуска команды.
Пример кода
CREATE TABLE binary12 (id INT PRIMARY KEY, str VARCHAR, xxx number) WITH "cache_name=binary12, value_type=PayloadTest15"Проверьте выполнение команды на исходном и принимающем кластерах. После проверки включите кеш в CDC.
Примечание
Чтобы использовать SQL-запросы в целевом кластере по данным, реплицированным CDC, установите одинаковый
VALUE_TYPEвlink:sql-reference/ddl#create-table[CREATE TABLE]для каждой таблицы в исходном и целевом кластерах.