Пример чтения данных из 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
-
Кластер ADS развернут согласно руководству Online-установка. Минимальная версия ADS — 4.0.0.b1.
-
Сервисы NiFi и Zookeeper установлены в кластере ADS.
ADB
-
Кластер ADB развернут согласно руководству Online-установка.
-
Пользователь с именем
my_user, привилегиямиSUPERUSERи паролем создан в БДadb. -
В БД
adbсоздана таблицаmy_tableи в нее добавлены несколько строк с данными.
Вход в систему под учетной записью 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:
$ 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)
-
Для выполнения подключения к Greengage DB создайте процессор GetGreengageRecord, откройте его конфигурацию и заполните параметры, связанные с таблицей Greengage DB.
Конфигурация процессора GetGreengageRecord
Конфигурация процессора GetGreengageRecord -
Перейдите в поле значения параметра Gpfdist Service, во всплывающем списке выберите Create new service… и в открывшемся окне создайте экземпляр сервиса StandardGpfdistService.
Создание экземпляра сервиса StandardGpfdistService
Создание экземпляра сервиса StandardGpfdistService -
После сохранения созданного экземпляра нажмите
. В открывшемся окне NiFi Flow Configuration → Controller Services откройте конфигурацию сервиса StandardGpfdistService и заполните необходимые параметры.
Конфигурация сервиса StandardGpfdistService
Конфигурация сервиса StandardGpfdistServiceПРИМЕЧАНИЕЗначение параметра Listening Port, заданное по умолчанию, может быть изменено, если данный порт уже используется.
-
Перейдите в поле значения параметра Database Connection Pooling Service, во всплывающем списке выберите Create new service… и в открывшемся окне создайте экземпляр сервиса DBCPConnectionPool для подключения к Greengage DB.
-
После сохранения созданного экземпляра нажмите
. В открывшемся окне NiFi Flow Configuration → Controller Services откройте конфигурацию сервиса DBCPConnectionPool и введите параметры базы данных Greengage DB и используемого драйвера.
Конфигурация сервиса GreengageDBCPConnectionPool
Конфигурация сервиса GreengageDBCPConnectionPool -
Закройте окно NiFi Flow Configuration → Controller Services и снова перейдите к конфигурации процессора GetGreengageRecord. В поле значения параметра Record Reader создайте экземпляр сервиса AvroRecordSetWriter для записи содержимого в бинарном формате Avro.
|
РЕКОМЕНДАЦИЯ
Если используется несколько процессоров GetGreengageRecord для переноса данных из разных таблиц одной базы Greengage DB, для значения Gpfdist Service выберите один и тот же созданный StandardGpfdistService. Создание экземпляров сервиса может быть выполнено до создания процессора. Это может быть использовано для подключения нескольких процессоров к одному сервису. Для создания сервиса выполните:
|
Подключение к серверу PostgreSQL (ADP)
-
Создайте процессор PutDatabaseRecord и откройте его конфигурацию. Этот процессор использует указанный RecordReader для ввода одной или нескольких записей из входящего потока файлов. Эти записи преобразуются в SQL-запросы для записи в таблицу PostgreSQL и выполняются как одна транзакция. Заполните параметры, связанные с используемой таблицей PostgreSQL.
Конфигурация процессора PutDatabaseRecord
Конфигурация процессора PutDatabaseRecord -
Перейдите в поле значения параметра Database Connection Pooling Service, во всплывающем списке выберите Create new service… и в открывшемся окне создайте экземпляр сервиса DBCPConnectionPool для подключения к серверу PostgreSQL.
-
После сохранения созданного экземпляра нажмите
, в открывшемся окне NiFi Flow Configuration → Controller Services откройте конфигурацию сервиса DBCPConnectionPool и введите параметры базы данных PostgreSQL и используемого драйвера.
Конфигурация сервиса PostgresDBCPConnectionPool
Конфигурация сервиса PostgresDBCPConnectionPool -
Закройте окно NiFi Flow Configuration → Controller Services и снова перейдите к конфигурации процессора PutDatabaseRecord. В поле значения параметра Record Reader создайте экземпляр сервиса AvroReader для чтения записей в формате Avro со встроенной схемой.
После создания все сервисы отображаются на странице NiFi Flow Configuration → Controller Services.
Сервис StandardGpfdistService отображается в статусе Invalid до запуска связанного с ним сервиса GreengageDBCPConnectionPool.
Запуск потока данных
Процессоры отображаются с ошибками из-за того, что сервисы, связанные с ними, не запущены.
Для запуска потока:
-
Поочередно запустите сервисы на странице NiFi Flow Configuration → Controller Services, кликнув на иконку
.
-
Запустите созданный поток данных.
Используя запросы к базе данных 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).