Механизм Change Data Capture (CDC) и межкластерная репликация#

Введение#

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.

cdc-design

Для persistence-кластеров используется следующий алгоритм работы механизма:

  1. При включенном механизме CDC серверный узел DataGrid создает в специальной директории db/cdc/{consistent_id} (значение по умолчанию) жесткую ссылку на каждый сегмент архива WAL-журнала.

  2. Затем приложение ignite-cdc.sh запускается на отдельной JVM и обрабатывает сегменты WAL-журнала, которые были перенесены в архив совсем недавно.

  3. После обработки сегмента приложением он удаляется. Само дисковое пространство освобождается при удалении обеих ссылок (в архиве WAL-журнала и в каталоге CDC).

CdcConsumer сохраняет состояние обработки в виде указателей на последние обработанные события в файлах cdc-*-state.bin. CdcConsumer может запрашивать сохранение (обновление) состояния у приложения ignite-cdc.sh. При запуске ignite-cdc.sh обработка событий будет продолжаться, начиная с последнего сохраненного состояния.

Внимание

Алгоритм работы механизма для in-memory-кластеров практически идентичен алгоритму работы CDC для persistence-кластеров. Отличие состоит в том, что в режиме in-memory узел записывает в WAL-журнал данные только тех областей данных (data region), для которых включен CDC.

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

Параметры конфигурации узла DataGrid#

Параметры конфигурации узла DataGrid#

Имя

Описание

Значение по умолчанию

DataRegionConfiguration#cdcEnabled

Свойство включает CDC на серверном узле для региона данных. Позволяет явно настроить регионы данных, для которых CDC включен

false

DataStorageConfiguration#cdcWalPath

Путь к каталогу CDC

db/wal/cdc

DataStorageConfiguration#walForceArchiveTimeout

Тайм-аут для принудительной отправки сегмента WAL-журнала в архив (даже если обработка сегмента не завершена).

Внимание

Необходимо обязательно установить тайм-аут, иначе репликация будет происходить только при архивировании WAL-сегмента по его заполненности, что непредсказуемо по длительности. Расхождения между кластерами в таком случае могут достигать нескольких минут или десятков минут (зависит от нагрузки). Если нагрузка низкая, ротация может происходить еще реже

-1 (disabled)

Имя

Описание

Значение по умолчанию

DataRegionConfiguration#cdcEnabled

Свойство включает CDC на серверном узле для региона данных. Позволяет явно настроить регионы данных, для которых CDC включен

false

DataStorageConfiguration#cdcWalPath

Путь к каталогу CDC

«db/wal/cdc»

DataStorageConfiguration#walForceArchiveTimeout

тайм-аут для принудительной отправки сегмента WAL-журнала в архив (даже если обработка сегмента не завершена).

Внимание: Необходимо обязательно установить тайм-аут, иначе репликация будет происходить только при архивировании WAL-сегмента по его заполненности, что непредсказуемо по длительности. Расхождения между кластерами в таком случае могут достигать нескольких минут или десятков минут (зависит от нагрузки). Если нагрузка низкая, ротация может происходить еще реже

-1 (disabled)

Параметры конфигурации приложения CDC#

CDC конфигурируется так же, как и узел DataGrid — через XML-файл Spring:

  • ignite-cdc.sh требует указания IgniteConfiguration и CdcConfiguration в одном конфигурационном файле;

  • IgniteConfiguration используется для установки общих параметров, например, пути к каталогу CDC, consistentId узла и других;

  • CdcConfiguration содержит непосредственно настройки приложения ignite-cdc.sh.

Имя

Описание

Значение по умолчанию

lockTimeout

тайм-аут для ожидания блокировки. При запуске CDC устанавливает блокировку на каталог во избежание параллельной обработки каталога CDC-журнала другим работающим приложением ignite-cdc.sh

1000 мс

checkFrequency

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

1000 мс

keepBinary

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

true

consumer

Реализация org.apache.ignite.cdc.CdcConsumer, обрабатывающая изменения в записях

null

metricExporterSpi

Массив SPI «экспортеров» для передачи метрик приложения CDC. Подробнее в разделе Конфигурация экспортеров и включение метрик в New Metrics System

null

Свойства кластера#

Свойства кластера, перечисленные в таблице ниже, позволяют настраивать CDC во время работы кластера (подробнее о работе со свойствами кластера см. на сайте Apache Ignite).

Имя

Описание

Значение по умолчанию

cdc.disabled

Принудительное отключение CDC во избежание переполнения диска. Полезно, если приложение CDC не работает в течение длительного времени. Внимание! Отключение CDC приведет к потере данных об изменениях

false

CDC API#

org.apache.ignite.cdc.CdcEvent

Полное описание интерфейса можно найти в официальной документации CdcEvent Interface.

Имя

Описание

key()

Ключ измененной записи

value()

Значение измененной записи. Метод вернет null при событии удаления записи

cacheId()

ID кеша, на котором происходит изменение. Значение метода равно значению параметра CACHE_ID из SYS.CACHES

partition()

Партиция DataGrid, в которой располагается измененная запись

primary()

Флаг основного (primary) узла. Возвращает true, если событие относится к основному узлу, false — к резервному (backup) узлу

version()

Comparable — версия измененной записи. DataGrid версионирует записи, что дает возможность упорядочить изменения каждой записи

org.apache.ignite.cdc.CdcConsumer

Полное описание интерфейса можно найти в официальной документации CdcConsumer Interface.

org.apache.ignite.cdc.CdcConsumer — обработчик (потребитель) событий-изменений. CdcConsumer может быть реализован пользователем самостоятельно. В дистрибутиве поставляется три реализации CDC (подробнее в разделе «Межкластерная репликация в DataGrid»).

Имя

Описание

void start(MetricRegistry)

Вызывается один раз при запуске приложения CDC. Реестр MetricRegistry должен использоваться для экспорта специфичных для потребителя метрик

boolean onEvents(Iterator<CdcEvent> events)

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

void stop()

Вызывается один раз при остановке приложения CDC

Метрики#

ignite-cdc.sh доступен тот же набор SPI для экспорта метрик, что и DataGrid. Следующие метрики предоставляются приложением (дополнительные метрики могут быть предоставлены потребителем):

Имя

Описание

CurrentSegmentIndex

Индекс текущего обрабатываемого сегмента WAL

CommittedSegmentIndex

Индекс сегмента WAL, содержащего последнее зафиксированное состояние

CommittedSegmentOffset

Зафиксированное смещение в байтах внутри сегмента WAL

LastSegmentConsumptionTime

Временная метка (в миллисекундах), указывающая на начало обработки последнего сегмента

BinaryMetaDir

Обрабатываемый приложением каталог с BinaryMeta

MarshallerDir

Обрабатываемый приложением каталог marshaller

CdcDir

Каталог CDC, из которого приложение считывает данные

SegmentConsumingTime

Время обработки WAL-сегмента в миллисекундах

Логирование#

ignite-cdc.sh использует ту же конфигурацию логирования, что и узел DataGrid. Единственное отличие — события записываются в файл ignite-cdc.log.

Жизненный цикл ignite-cdc.sh#

Внимание

В приложении ignite-cdc.sh реализован подход fail-fast. В случае любой ошибки приложение прекращает работу. Процедура перезапуска должна быть настроена средствами операционной системы.

После запуска приложение ignite-cdc.sh:

  1. Находит необходимые общие каталоги. Используются значения из предоставленной IgniteConfiguration.

  2. Блокирует CDC каталог.

  3. Загружает сохраненное состояние («отметку» последнего обработанного события).

  4. Запускает CdcConsumer.

  5. Запускает бесконечный цикл ожидания и обработки новых доступных WAL-сегментов.

  6. Останавливает потребитель данных в случае ошибки или получения сигнала об остановке.

Обработка пропущенных сегментов#

Механизм 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).

В зависимости от архитектуры сети и требований безопасности доступны три схемы межкластерной репликации:

  1. ⁣DataGrid to DataGrid через тонкий клиент Java (Java thin client).

  2. ⁣DataGrid to DataGrid через клиентский узел.

  3. ⁣DataGrid с использованием Platform V Corax.

Примечание

Все три схемы CDC поддерживают репликацию изменений TypeMapping и BinaryType.

Для разрешения возможных конфликтов изменения данных при использовании Active-Active режима репликации необходимо настроить модуль ConflictResolver для каждой схемы межкластерной репликации. Подробное описание этого модуля содержится в следующем разделе.

Репликация DataGrid to DataGrid через тонкий клиент Java (Java thin client)#

IgniteToIgniteCdcStreamer запускает тонкий клиент Java, подключающийся к удаленному кластеру-получателю. При наличии подключения к удаленному кластеру тонкий клиент Java будет реплицировать в него изменения данных, захваченные на кластере-источнике:

Ignite2KafkaJavathin

Внимание

Экземпляры приложения ignite-cdc.sh должны быть настроены и запущены на каждом серверном узле кластера-источника для передачи всех изменений данных.

Данная схема репликации требует возможности организации прямого соединения между двумя кластерами DataGrid.

Примечание

На текущий момент данная схема позволяет обеспечить репликацию только между двумя кластерами. Для соединения более двух кластеров можно воспользоваться схемой из раздела «Репликация DataGrid с использованием Platform V Corax».

Описание конфигурации#

Имя

Описание

Значение по умолчанию

caches

Набор (set) имен реплицируемых кешей

null

destinationClientConfiguration

Конфигурация клиентского узла, реплицирующего изменения на кластер-получатель

null

onlyPrimary

Свойство отвечает за то, какие партиции будут реплицироваться: либо только основные (primary) партиции, либо все, чтобы в случае выхода основной партиции из строя изменения поступили от резервных (backup) партиций

false

maxBatchSize

Максимальный размер пакета изменений, которые необходимо отправить

1024

Метрики#

Имя

Описание

EventsCount

Количество событий, примененных к кластеру-получателю

LastEventTime

Отметка времени последнего примененного события

TypesCount

Количество обработанных изменений в BinaryMeta

MappingsCount

Количество обработанных изменений в Marshaller

Репликация DataGrid to DataGrid через клиентский узел#

IgniteToIgniteCdcStreamer запускает клиентский узел, подключающийся к удаленному кластеру-получателю. При наличии подключения к удаленному кластеру клиентский узел будет реплицировать в него изменения данных, захваченные на кластере-источнике.

IgnitetoIgnite

Внимание

Экземпляры приложения ignite-cdc.sh должны быть настроены и запущены на каждом серверном узле кластера-источника для передачи всех изменений данных.

Данная схема репликации требует возможности организации прямого соединения между двумя кластерами DataGrid.

Примечание

На текущий момент данная схема позволяет обеспечить репликацию только между двумя кластерами. Для соединения более двух кластеров можно воспользоваться схемой из раздела «Репликация DataGrid с использованием Platform V Corax».

Описание конфигурации#

Имя

Описание

Значение по умолчанию

caches

Набор (set) имен реплицируемых кешей

null

destinationIgniteConfiguration

Конфигурация клиентского узла, реплицирующего изменения на кластер-получатель

null

onlyPrimary

Свойство отвечает за то, какие партиции будут реплицироваться: либо только основные (primary) партиции, либо все, чтобы в случае выхода основной партиции из строя изменения поступили от резервных (backup) партиций

false

maxBatchSize

Максимальный размер пакета изменений, которые необходимо отправить

1024

Метрики#

Имя

Описание

EventsCount

Количество событий, реплицированных на кластер-получатель

LastEventTime

Отметка времени последнего примененного события

TypesCount

Количество обработанных изменений в BinaryMeta

MappingsCount

Количество обработанных изменений в Marshaller

Репликация DataGrid с использованием Platform V Corax#

Этот способ репликации изменений между кластерами требует настройки двух приложений:

  1. ignite-cdc.sh с классом org.apache.ignite.cdc.kafka.IgniteToKafkaCdcStreamer, который будет захватывать изменения из кластера-источника и записывать их в топик Corax.

  2. 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). Репликация возможна как между кластерами, содержащимися в одном ЦОД, так и между кластерами, содержащимися в разных ЦОД. Различий в схеме репликации и ее настройке нет.

Ignite2Kafka

IgniteToKafkaCdcStreamer отправляет захваченные на кластере-источнике изменения в топик Platform V Corax. KafkaToIgniteCdcStreamerвычитывает изменения из топик и применяет на удаленном сервере-приемнике.

Для двусторонней репликации (Active-Active) запись событий CDC и IgniteToKafkaCdcStreamer (ignite-cdc.sh) включаются на всех кластерах. Для применения изменений на каждом приемнике необходимо запустить как минимум по одному экземпляру KafkaToIgniteCdcStreamer для каждого из топиков с событиями CDC других кластеров-источников (топик кластера-приемника вычитывать на самом кластре-приемнике не нужно). Кроме того, для разрешения конфликтных изменений необходмо настроить CacheVersionConflictResolver при помощи соответствующего плагина.

Описание конфигурации#

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

Имя

Описание

Значение по умолчанию

caches

Набор (set) имен реплицируемых кешей

null

kafkaProperties

Свойства поставщика (producer) Corax

null

topic

Имя топика Corax для событий CDC

null

kafkaParts

Количество партиций топика Corax для событий CDC

null

metadataTopic

Топик для репликации изменений TypeMapping и BinaryType

null

onlyPrimary

Свойство отвечает за то, какие партиции будут реплицироваться: либо только основные (primary) партиции, либо все, чтобы в случае выхода основной партиции из строя изменения поступили от резервных (backup) партиций

false

maxBatchSize

Максимальный размер конкурентно создаваемых записей Corax.

По достижении этого значения CdcStreamer ожидает подтверждения получения от Corax, а затем выполняет коммит смещения CDC

1024

kafkaRequestTimeout

тайм-аут запроса Corax, мс

3000

Параметр 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 является обязательным.

Метрики#

Имя

Описание

EventsCount

Количество событий, примененных к Corax

LastEventTime

Отметка времени последнего примененного события

TypesCount

Количество обработанных изменений в BinaryMeta

MappingsCount

Количество обработанных изменений в Marshaller

BytesSent

Количество байтов, отправленных в Corax

MarkersCount

Количество маркеров об обработанных изменениях в Marshaller и BinaryMeta, отправленных в 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.

Имя

Описание

Значение по умолчанию

caches

Набор (set) имен реплицируемых кешей

null

kafkaConsumerPollTimeout

Тайм-аут запроса #poll() Corax (мс)

3000

kafkaPartsFrom

Начальная партиция Corax (включительно) для топика событий CDC

-1

kafkaPartsTo

Конечная партиция Corax (не включая) для топика CDC событий

-1

kafkaRequestTimeout

Тайм-аут запроса Corax (мс)

3000

maxBatchSize

Максимальный размер конкурентно создаваемых записей Corax.

По достижении этого значения CdcStreamer ожидает подтверждения получения от Corax, а затем выполняет коммит смещения CDC

1024

metadataConsumerGroup

Группа для KafkaConsumer, который читает данные из топика метаданных

ignite-metadata-update-<kafkaPartsFrom>-<kafkaPartsTo>

metadataTopic

Топик для репликации изменений TypeMapping и BinaryType

null

threadCount

Количество потоков для выполнения консьюмерами.

Каждый поток опрашивает свой диапазон партиций в цикле

16

topic

Имя топика 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 является обязательным.

Метрики#

Имя

Описание

EventsReceivedCount

Количество событий, полученных из Corax

LastEventReceivedTime

Отметка времени последнего события, полученного из Corax

EventsSentCount

Количество событий, отправленных в кластер-приемник

LastBatchSentTime

Отметка времени последнего события, примененного в кластере-приемнике

MarkersCount

Количество маркеров об обработанных изменениях в Marshaller и BinaryMeta, полученных из 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#

Имя

Описание

Значение по умолчанию

clusterId

Локальный идентификатор кластера. Может принимать значения от 1 до 31

null

caches

Набор (set) кешей для обработки плагином

null

conflictResolveField

Опциональное поле для разрешения конфликта.
Поле должно реализовывать интерфейс java.lang.Comparable. Значение должно задаваться в объектах, которые прикладной код записывает в реплицируемые кеши DataGrid

null

conflictResolver

Пользовательское средство разрешения конфликта. Опционально.
Поле должно реализовывать интерфейс CacheVersionConflictResolver

null

Механизм работы модуля#

Реплицированные данные содержат дополнительную информацию, в частности, версии записей (CacheEntryVersion) из кластера-источника. Алгоритм разрешения конфликтов по умолчанию основан на сравнении версий записи и поля conflictResolveField.

Разрешение конфликтов на основе версии записи#

Этот подход обеспечивает окончательную гарантию согласованности только в том случае, если каждая запись обновляется только из одного кластера.

Алгоритм разрешения конфликтов:

  1. Изменение из «локального» кластера (clusterId соответствует текущему кластеру) всегда применяется безусловно. Любые реплицированные данные могут быть перезаписаны локально.

  2. Если имеющаяся запись и запись-кандидат из одного и того же кластера (clusterId совпадает), то будет выбрано значение с бо́льшим значением CacheEntryVersion.

  3. При невозможности разрешить конфликт изменение не принимается, что приводит к потере консистентности между кластерами. Ошибка разрешения конфликта выводится в log-файл.

Внимание

При таком подходе не происходит репликация обновлений и удалений записей из кластера-получателя на кластер-источник.

Для обеспечения максимальной производительности и во избежание ситуации, описанной в пункте 3, рекомендуется избегать пересечений по ключам при проектировании, а при наличии пересечений обоснованно выбирать поля для сравнения изменений.

Разрешение конфликтов на основе значения записи#

Этот подход обеспечивает конечную гарантию согласованности, даже если запись была обновлена из разных кластеров.

Примечание

Поле разрешения конфликта, указанное в conflictResolveField, должно содержать монотонно возрастающее значение, предоставленное пользователем. Например, идентификатор запроса или метку времени.

Внимание

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

Алгоритм разрешения конфликтов:

  1. Изменение из «локального» кластера (clusterId соответствует текущему кластеру) всегда применяется безусловно. Любые реплицированные данные могут быть перезаписаны локально.

  2. Если имеющаяся запись и запись-кандидат из одного и того же кластера (clusterId совпадает), то будет выбрано значение с бо́льшим значением CacheEntryVersion.

  3. Если конфликт еще не разрешен и задано поле conflictResolveField, то оно сравнивается у текущей записи и у записи-кандидата — выбирается запись с бо́льшим значением conflictResolveField.

  4. При невозможности разрешить конфликт изменение не принимается, что приводит к потере консистентности между кластерами. Ошибка разрешения конфликта выводится в log-файл.

Пользовательские правила для разрешения конфликтов#

Пользователь может задать собственные правила для разрешения конфликтов, исходя из характера данных и операций. Это может быть полезно, когда стандартные стратегии для разрешения конфликтов неприменимы.

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

Пользовательское средство разрешения конфликтов может быть задано с помощью conflictResolver и позволяет сравнивать или объединять конфликтные данные любым требуемым способом.

Пример настройки плагина CacheVersionConflictResolver#

Указывается в serverExampleConfig.xml.

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.

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.

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 на каждом серверном узле. Для запуска выполните следующие действия:

  1. Выполните команду ignite-cdc.sh ignite-1984.xml.

  2. Выполните команду ignite-cdc.sh ignite-2029.xml.

DataGrid to DataGrid через тонкий клиент Java#

Примечание

Конфигурация ниже предназначена для локального развертывания (на localhost). Настройки указаны для примера.

Настройка аналогична настройке «Репликация DataGrid to DataGrid через клиентский узел». Отличие в том, что в экземпляре CdcConfiguration в поле consumer необходимо указать IgniteToIgniteClientStreamer.

XML#
<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.

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.

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.

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.

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.

XML#
<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:

  1. Запуск DataGrid to Platform V Corax:

    ignite-cdc.sh ignite-dc1-kafka.xml
    
    ignite-cdc.sh ignite-dc2-kafka.xml
    
  2. Запуск Platform V Corax to DataGrid:

    kafka-to-ignite.sh kafka2ignite_dc1.xml
    
    kafka-to-ignite.sh kafka2ignite_dc2.xml
    

Возможные проблемы в работе репликации и методы их устранения#

  1. Чтобы CdcStreamer не завершил работу с ошибкой, до запуска репликации необходимо:

    1. Активировать кластер-получатель.

    2. Создать реплицируемый кеш на кластере-получателе.

  2. Во избежание возникновения сбоев на серверных узлах cdcWalPath и walArchivePath должны быть расположены в одном разделе файловой системы. В противном случае серверный узел завершит работу с ошибкой.

  3. Во избежание зацикливания сообщений при репликации в режиме Active-Active необходимо убедиться, что:

    1. Для всех реплицируемых кешей настроен conflictResolver.

    2. В conflictResolver указаны разные clusterId для разных кластеров.

    3. Если conflictResolveField не указан в conflictResolver, то для уже имеющихся локально записей конфликтные изменения с удаленного кластера разрешаться не будут.

  4. Во избежание ошибок при работе с одинаково названными 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] для каждой таблицы в исходном и целевом кластерах.