Пример чтения данных из ADB с помощью NiFi ADB Connector

Обзор

Для иллюстрации работы NiFi ADB Connector в статье показана реализация чтения данных из таблицы ADB (на основе Greengage DB) и загрузки данных в таблицу базы данных PostgreSQL.

Создание NiFi ADB Connector выполняется в пользовательском интерфейсе NiFi. Функция чтения данных из ADB доступна начиная с ADS 4.0.0.b1.

Предварительные требования

Ниже описано окружение, используемое для создания NiFi ADB Connector.

ADS

При настройке сервисов DBCPConnectionPool для подключения к серверам PostgreSQL (ADP) и Greengage DB (ADB) используются значения, связанные с конфигурацией кластеров ADS, ADP и ADB:

  • Database Connection URL — ссылка в формате jdbc:postgresql://<host>:5432/<database>, где:

    • <host> — хост, на котором установлен компонент ADPG (для подключения к серверу PostgreSQL) или компонент ADB Master (для подключения к серверу Greengage DB);

    • <database> — наименование базы данных, в которой создана используемая таблица.

    Пример ссылки: jdbc:postgresql://10.92.38.119:5432/adb.

  • Database Driver Class Name — имя класса драйвера JDBC PostgreSQL, который позволяет программам подключаться к базе данных PostgreSQL, используя стандартный код Java (например, org.postgresql.Driver).

  • Database Driver Location(s) — место расположения файла драйвера в формате file:<path><driver_name>, где:

    • <path> — путь к двоичному JAR-файлу драйвера, размещенному на хостах, где установлен NiFi.

    • <driver_name> — имя JAR-файла драйвера, соответствующего используемой версии Greengage DB или PostgreSQL.

    Пример указания пути к файлу драйвера, соответствующего ADB 6.30.0: file:/tmp/postgresql-42.2.27.jar.

    Команда, приведенная ниже, позволяет загрузить необходимую версию JAR-файла в нужную директорию:

    $ wget https://jdbc.postgresql.org/download/postgresql-42.2.27.jar

ADB

  • Кластер ADB развернут согласно руководству Online-установка.

  • Пользователь с именем my_user, привилегиями SUPERUSER и паролем создан в БД adb.

  • В БД adb создана таблица my_table и в нее добавлены несколько строк с данными.

Пример настройки БД adb

Вход в систему под учетной записью gpadmin:

$ sudo su - gpadmin

Подключение к БД через psql:

$ psql adb

Создание пользователя с ролью SUPERUSER:

CREATE USER my_user WITH SUPERUSER PASSWORD 'P@ssword';
ВНИМАНИЕ
Роль пользователя SUPERUSER используется только для тестовых целей.

Создание тестовой таблицы:

CREATE TABLE my_table (
    id BIGSERIAL PRIMARY KEY,
    name VARCHAR(100) NOT NULL,
    country VARCHAR(100),
    updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);

Заполнение таблицы данными (приведено как пример, для наглядности необходимо несколько строк):

INSERT INTO my_table(name, country) VALUES ('John Jones','USA');
  • Файл pg_hba.conf настроен для обеспечения доступа пользователя с хоста, на котором установлен сервис NiFi кластера ADS. Для этого в поле Custom pg_hba section на странице конфигурационных параметров сервиса ADB добавлена запись об адресе хоста:

    host      all             all        10.92.38.11/24       trust
  • В Interconnect-сети, к которой подключены хосты кластера, установлен MTU=9000(jumbo frame), чтобы пакеты, формируемые ADB (gp_max_packet_size + overhead), помещались в эти фреймы целиком. Для получения более подробной информации о требованиях к сети кластера ADB обратитесь к статье Требования к сети.

ПРИМЕЧАНИЕ

Для получения информации о работе в ADB обратитесь к статьям:

ADP

  • Кластер ADP развернут согласно руководству Online-установка.

  • Пользователь с именем my_user, привилегиями SUPERUSER и паролем создан в БД postgres.

  • В БД postgres создана таблица my_table, тип и наименование столбцов которой совпадают с копируемой таблицей из ADB.

Пример настройки БД postgres

Подключение к базе postgres:

$ sudo su - postgres
$ psql

Создание пользователя с ролью SUPERUSER:

CREATE USER my_user WITH SUPERUSER PASSWORD 'P@ssword';
ВНИМАНИЕ
Роль пользователя SUPERUSER используется только для тестовых целей.

Создание тестовой таблицы:

CREATE TABLE my_table (
    id BIGSERIAL PRIMARY KEY,
    name VARCHAR(100) NOT NULL,
    country VARCHAR(100),
    updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
  • Файл pg_hba.conf настроен для обеспечения доступа пользователя с хоста, на котором установлен сервис NiFi кластера ADS. Для этого в поле PG_HBA на странице конфигурационных параметров сервиса ADPG добавлена запись об адресе хоста и пользователе:

    host      all             all        10.92.38.11/24       trust
ПРИМЕЧАНИЕ

Для получения информации о работе в ADP обратитесь к статьям:

Подключение к Greengage DB (ADB)

  1. Для выполнения подключения к Greengage DB создайте процессор GetGreengageRecord, откройте его конфигурацию и заполните параметры, связанные с таблицей Greengage DB.

    Конфигурация процессора GetGreengageRecord
    Конфигурация процессора GetGreengageRecord
    Конфигурация процессора GetGreengageRecord
    Конфигурация процессора GetGreengageRecord
  2. Перейдите в поле значения параметра Gpfdist Service, во всплывающем списке выберите Create new service…​ и в открывшемся окне создайте экземпляр сервиса StandardGpfdistService.

    Создание экземпляра сервиса StandardGpfdistService
    Создание экземпляра сервиса StandardGpfdistService
    Создание экземпляра сервиса StandardGpfdistService
    Создание экземпляра сервиса StandardGpfdistService
  3. После сохранения созданного экземпляра нажмите arrow2 light. В открывшемся окне NiFi Flow Configuration → Controller Services откройте конфигурацию сервиса StandardGpfdistService и заполните необходимые параметры.

    Конфигурация сервиса StandardGpfdistService
    Конфигурация сервиса StandardGpfdistService
    Конфигурация сервиса StandardGpfdistService
    Конфигурация сервиса StandardGpfdistService
    ПРИМЕЧАНИЕ

    Значение параметра Listening Port, заданное по умолчанию, может быть изменено, если данный порт уже используется.

  4. Перейдите в поле значения параметра Database Connection Pooling Service, во всплывающем списке выберите Create new service…​ и в открывшемся окне создайте экземпляр сервиса DBCPConnectionPool для подключения к Greengage DB.

  5. После сохранения созданного экземпляра нажмите arrow2 light. В открывшемся окне NiFi Flow Configuration → Controller Services откройте конфигурацию сервиса DBCPConnectionPool и введите параметры базы данных Greengage DB и используемого драйвера.

    Конфигурация сервиса GreengageDBCPConnectionPool
    Конфигурация сервиса GreengageDBCPConnectionPool
    Конфигурация сервиса GreengageDBCPConnectionPool
    Конфигурация сервиса GreengageDBCPConnectionPool
  6. Закройте окно NiFi Flow Configuration → Controller Services и снова перейдите к конфигурации процессора GetGreengageRecord. В поле значения параметра Record Reader создайте экземпляр сервиса AvroRecordSetWriter для записи содержимого в бинарном формате Avro.

РЕКОМЕНДАЦИЯ

Если используется несколько процессоров GetGreengageRecord для переноса данных из разных таблиц одной базы Greengage DB, для значения Gpfdist Service выберите один и тот же созданный StandardGpfdistService.

Создание экземпляров сервиса может быть выполнено до создания процессора. Это может быть использовано для подключения нескольких процессоров к одному сервису. Для создания сервиса выполните:

  1. Правой кнопкой мыши кликните в пустом поле потока и выберите Configure в открывшемся контекстном меню.

  2. В окне NiFi Flow Configuration перейдите на вкладку Controller Services и кликните на +.

  3. Выберите нужный сервис из списка и создайте экземпляр сервиса в открывшемся окне.

Подключение к серверу PostgreSQL (ADP)

  1. Создайте процессор PutDatabaseRecord и откройте его конфигурацию. Этот процессор использует указанный RecordReader для ввода одной или нескольких записей из входящего потока файлов. Эти записи преобразуются в SQL-запросы для записи в таблицу PostgreSQL и выполняются как одна транзакция. Заполните параметры, связанные с используемой таблицей PostgreSQL.

    Конфигурация процессора PutDatabaseRecord
    Конфигурация процессора PutDatabaseRecord
    Конфигурация процессора PutDatabaseRecord
    Конфигурация процессора PutDatabaseRecord
  2. Перейдите в поле значения параметра Database Connection Pooling Service, во всплывающем списке выберите Create new service…​ и в открывшемся окне создайте экземпляр сервиса DBCPConnectionPool для подключения к серверу PostgreSQL.

  3. После сохранения созданного экземпляра нажмите arrow2 light, в открывшемся окне NiFi Flow Configuration → Controller Services откройте конфигурацию сервиса DBCPConnectionPool и введите параметры базы данных PostgreSQL и используемого драйвера.

    Конфигурация сервиса PostgresDBCPConnectionPool
    Конфигурация сервиса PostgresDBCPConnectionPool
    Конфигурация сервиса PostgressDBCPConnectionPool
    Конфигурация сервиса PostgresDBCPConnectionPool
  4. Закройте окно NiFi Flow Configuration → Controller Services и снова перейдите к конфигурации процессора PutDatabaseRecord. В поле значения параметра Record Reader создайте экземпляр сервиса AvroReader для чтения записей в формате Avro со встроенной схемой.

После создания все сервисы отображаются на странице NiFi Flow Configuration → Controller Services.

Созданные сервисы
Созданные сервисы
Созданные сервисы
Созданные сервисы

Сервис StandardGpfdistService отображается в статусе Invalid до запуска связанного с ним сервиса GreengageDBCPConnectionPool.

Запуск потока данных

Создайте и настройте подключение между процессорами.

Созданные и соединенные процессоры
Созданные и соединенные процессоры
Созданные и соединенные процессоры
Созданные и соединенные процессоры

Процессоры отображаются с ошибками из-за того, что сервисы, связанные с ними, не запущены.

Для запуска потока:

  1. Поочередно запустите сервисы на странице NiFi Flow Configuration → Controller Services, кликнув на иконку nifi ui oper 02.

  2. Запустите созданный поток данных.

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

SELECT * FROM my_table;

Обновление потока данных

В пайплайне, описанном в этой статье, для получения новых строк из базы adb требуется выполнить очистку состояния (state) кластера. Это вызовет повторное чтение таблицы целиком. В базу postgres будут записаны новые строки, а имеющиеся обновлены (если в них были внесены изменения).

Для очистки состояния кластера правой кнопкой мыши вызовите контекстное меню процессора, остановите процессор и выберите команду View state. На открывшейся вкладке Component State кликните на Clear state, затем заново запустите процессор.

Очистка состояния кластера
Очистка состояния кластера
Очистка состояния кластера
Очистка состояния кластера

Вкладка Component State отображает состояние кластера — список снепшотов, сохраняющих для процессора максимальные наблюдаемые значения по каждому исполнителю (worker) и столбцу для инкрементальной (пошаговой) выгрузки из Greengage DB.

Исполнитель — одна из параллельных задач процессора, отвечающих за выгрузку данных из сегментов. Количество исполнителей определяется как минимальное значение между количеством сегментов Greengage и произведением общего числа узлов NiFi на значение параметра Node Parallel Factor.

Режим инкрементальной выгрузки

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

Режим пошаговой выгрузки поддерживается для следующих типов данных: smallint, integer, bigint, real, double precision, numeric, date, time, timestamp, timestamp with time zone.

Работа пошаговой выгрузки

Для включения режима пошаговой выгрузки укажите имя одного или нескольких столбцов в поле параметра Maximum-value Columns Names.

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

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

Если нового диапазона нет, выгрузка для этого исполнителя пропускается.

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

При перезапуске процессора инкрементальная позиция восстанавливается из состояния кластера.

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

После смены значений параметра Maximum-value Columns Names (списка столбцов) выполните очистку состояния кластера.

Работа без пошаговой выгрузки

Если значение параметра Maximum-value Columns Names не установлено (например, как в пайплайне, описанном в этой статье), процессор работает в режиме полной загрузки (full load):

  • каждый worker выполняет полную выгрузку один раз;

  • процессор сохраняет последний снепшот исходной таблицы в состоянии кластера для каждого исполнителя;

  • следующие циклы пропускаются, пока состояние не будет очищено.

Идемпотентный режим

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

Строки, вставленные в таблицу Greengage DB (ADB) во время выгрузки, гарантированно попадут в следующий цикл.

State полностью записывается (коммитится) только после полного успеха цикла: при любой ошибке верхняя граница не изменяется и диапазон перечитывается заново, поэтому возможны дубликаты строк.

Изменение значения Node Parallel Factor меняет распределение исполнителей между сегментами: для новых исполнителей состояние пусто и они выполнят первичную выгрузку всего диапазона (возможны дубли с уже выгруженным).

Исходя из этого, процессор GetGreengageRecord поддерживает семантику доставки at least once (сообщения будут обработаны хотя бы один раз).

Для организации семантики exactly once (каждое сообщение будет обработано один и только один раз) на стороне записи нужен идемпотентный режим для исключения дублирования строк:

  • UPSERT/UPDATE по ключу (например, в пайплайне GetGreengageRecord → PutDatabaseRecord);

  • дедупликация по reading_id в конечном хранилище (например, в пайплайне GetGreengageRecord → UpdateRecord → PutGreengageRecord).

Нашли ошибку? Выделите текст и нажмите Ctrl+Enter чтобы сообщить о ней