Провайдер HBase
Обзор
Провайдер HBase позволяет DAG взаимодействовать с HBase. Он предоставляет хуки (hook), операторы (operator) и сенсоры (sensor) для выполнения административных задач и операций с данными в таблицах HBase, а также резервного копирования и восстановления.
Провайдер HBase позволяет:
-
создавать и удалять таблицы;
-
выполнять операции вставки, обновления, чтения и удаления строк;
-
сканировать таблицы;
-
выполнять пакетные операции;
-
создавать резервные копии и использовать их для восстановления;
-
управлять наборами резервного копирования;
-
проверять доступность таблиц и строк.
Провайдер соответствует стандартной архитектуре провайдеров в Airflow, в которой операторы выполняют задачи в DAG, хуки обеспечивают низкоуровневое взаимодействие с внешними системами, а сенсоры ожидают выполнения условий.
В отличие от универсальных провайдеров Airflow, провайдер HBase поддерживает два механизма работы:
-
Thrift2 API для операций с данными в режиме онлайн;
-
утилиты HBase CLI для операций резервного копирования и восстановления.
Архитектура
Провайдер состоит из нескольких модулей Python:
-
Hooks
Этот модуль отвечает за взаимодействие с HBase. Он содержит классы, устанавливающие соединения и предоставляющие интерфейс Python для выполнения операций HBase. В провайдере предусмотрены две реализации хуков:
-
HBaseThriftHook— взаимодействует с компонентом HBase Thrift2 Server и отвечает за управление таблицами, обработку данных и операции с метаданными. -
HBaseCLIHook— выполняет команды клиента HBase и используется в рабочих процессах резервного копирования и восстановления.
-
-
Operators
Этот модуль содержит задачи Airflow, представляющие отдельные операции в HBase. Каждый оператор проверяет параметры, инициализирует соответствующий хук, вызывает необходимый метод и возвращает результат выполнения в Airflow.
-
Sensors
Модуль предоставляет сенсоры Airflow, которые периодически проверяют состояние ресурсов HBase. В отличие от операторов, сенсоры не изменяют данные. Они многократно выполняют небольшие запросы к HBase, проверяя выполнение указанного условия в течение заданного времени (тайм-аута). Например, ожидают создания таблицы перед загрузкой данных или появления определенной строки в результате выполнения другого рабочего процесса.
-
Utils
Модуль содержит вспомогательные классы, используемые несколькими компонентами провайдера. К ним относятся: перечисления (enumerations), описывающие типы резервного копирования и поведение операторов, константы, используемые в провайдере, механизмы повторных попыток и другие вспомогательные функции, поддерживающие реализацию.
Эти модули работают как единый пайплайн. Во время выполнения DAG Airflow создает экземпляр оператора, который выбирает соответствующий хук в зависимости от требуемой операции. Хук устанавливает соединение с HBase по протоколу Thrift2 или через HBase CLI, выполняет требуемую операцию, при необходимости преобразует ответ в объект Python и возвращает результат оператору. Затем оператор передает статус выполнения задачи планировщику Airflow.
Требования
Требования зависят от используемой функциональности.
Если у вас уже есть кластер ADH с необходимыми компонентами, вы можете задействовать его совместно с кластером ADO, используя общие узлы.
Операции с данными
Для работы операторов и сенсоров необходимы:
-
HBase Thrift2 Server, запущенный в кластере ADH и доступный по сети для воркеров Airflow;
-
настроенное подключение HBase в Airflow.
Конфигурация
Провайдер HBase использует тип подключения Airflow hbase для хранения параметров, необходимых для взаимодействия с HBase. Операторы и сенсоры принимают параметр hbase_conn_id, который ссылается на настроенное подключение Airflow. Если идентификатор подключения не указан явно, по умолчанию используется подключение hbase_thrift2.
Провайдер поддерживает два типа подключений:
-
hbase— нативный тип подключения, реализованный провайдером. Он рекомендуется для всех сценариев работы с HBase, поскольку поддерживает полный набор специфичных для провайдера параметров конфигурации, включая политики повторных попыток, пул соединений, параметры SSL/TLS и конфигурацию CLI. -
generic— может использоваться для подключения к произвольным серверам Thrift. Этот тип подключения предназначен главным образом для сценариев совместимости и не предоставляет специфичные для HBase параметры конфигурации, поддерживаемые нативным типом подключения.
Стандартная конфигурация
Ниже приведен список стандартных полей подключения в Airflow, используемых провайдером. Остальная конфигурация задается в поле Extra в формате JSON.
| Поле | Описание |
|---|---|
Connection type |
Должно быть установлено значение |
Host |
Указывает имя хоста или IP-адрес HBase Thrift2 Server |
Port |
Указывает порт для Thrift2. Если не задан, используется стандартный порт HBase Thrift2 Server ( |
Extra |
Содержит специфичные для провайдера параметры конфигурации |
Пример подключения Thrift2 без аутентификации:
-
Connection Type: hbase -
Host: hbase-server.example.com -
Port: 9090 -
Extra:{ "use_http": false }
Дополнительные параметры
Следующие параметры управляют поведением клиента Thrift2.
| Параметр | Описание | Значение по умолчанию |
|---|---|---|
timeout |
Определяет тайм-аут запроса Thrift в миллисекундах |
30000 |
namespace |
Указывает пространство имен HBase, используемое по умолчанию операторами, которые не определяют его явно |
default |
use_http |
Включает транспорт HTTP вместо используемого по умолчанию бинарного транспорта через сокет |
false |
retry_max_attempts |
Определяет максимальное количество повторных попыток после сбоя соединения |
3 |
retry_delay |
Определяет начальную задержку между повторными попытками в секундах |
1.0 |
retry_backoff_factor |
Определяет множитель, используемый для увеличения задержки после каждой неудачной повторной попытки |
2.0 |
Пул соединений
По умолчанию провайдер использует одно соединение Thrift для каждого выполнения задачи. Эта стратегия, реализованная в Thrift2Strategy, подходит для большинства рабочих нагрузок и минимизирует потребление ресурсов.
Для производственных сред, в которых выполняются большие пакетные операции или несколько задач HBase одновременно, провайдер также поддерживает использование пула соединений через PooledThrift2Strategy. Вместо создания нового соединения для каждого запроса провайдер поддерживает пул повторно используемых соединений Thrift, которые могут совместно использоваться различными операциями. Повторное использование уже установленных соединений снижает накладные расходы на их создание и значительно увеличивает пропускную способность при интенсивных рабочих нагрузках.
По умолчанию пул соединений отключен. Чтобы включить его, добавьте раздел connection_pool в поле подключения Extra.
{
"connection_pool": {
"enabled": true,
"size": 10,
"timeout": 30
}
}
Доступные параметры пула соединений приведены ниже.
| Параметр | Описание | Значение по умолчанию |
|---|---|---|
enabled |
Включает пул соединений |
false |
size |
Указывает максимальное количество соединений Thrift, поддерживаемых в пуле |
10 |
timeout |
Определяет максимальное время в секундах, в течение которого задача ожидает доступное соединение перед возвратом ошибки |
30 |
|
РЕКОМЕНДАЦИЯ
Для производственных сред, выполняющих пакетную обработку или несколько задач HBase одновременно, рекомендуется использовать пул соединений.
|
SSL/TLS
Провайдер поддерживает защищенное взаимодействие с компонентом HBase Thrift2 Server. Параметры SSL/TLS задаются в поле подключения Extra.
| Параметр | Описание | Значение по умолчанию |
|---|---|---|
ca_certs |
Указывает путь к набору доверенных сертификатов центра сертификации (CA) |
— |
validate |
Определяет, выполняется ли проверка сертификата сервера при установлении соединения |
true |
use_http |
Должен быть установлен в значение |
false |
Следующий пример включает SSL-соединение с проверкой сертификата сервера:
{
"ca_certs": "/etc/ssl/hbase_certs.pem",
"validate": true,
"use_http": true
}
Параметры CLI
Операторы резервного копирования и восстановления используют HBase Backup CLI на хостах Airflow вместо взаимодействия с HBase через Thrift2. Ниже приведены опциональные параметры для настройки среды выполнения.
| Параметр | Описание | Значение по умолчанию |
|---|---|---|
java_home |
Указывает установку Java, используемую для выполнения команд HBase CLI |
/usr/lib/jvm/java-arenadata-openjdk-8 |
hbase_home |
Указывает директорию установки HBase, содержащую клиентские утилиты |
/usr/lib/hbase |
В средах, где компонент HBase Client установлен из ADH с использованием общих хостов, эти параметры обычно не требуют изменения, поскольку уже указывают на стандартные директории установки в ADH.
Хуки
Хуки реализуют уровень взаимодействия провайдера. Провайдер содержит две реализации хуков, поскольку сам HBase предоставляет разные интерфейсы для различных категорий операций.
HBaseThriftHook
HBaseThriftHook — основной интерфейс взаимодействия, используемый провайдером. Он устанавливает соединение с HBase Thrift2 Server и предоставляет методы для администрирования таблиц, обработки данных, получения метаданных и пакетной обработки. Все операторы для работы с данными и сенсоры используют этот хук для взаимодействия с кластером HBase.
Хук предоставляет методы для наиболее распространенных операций HBase, включая управление таблицами, обработку строк, сканирование таблиц и пакетную обработку.
| Метод | Описание |
|---|---|
table_exists() |
Проверяет существование указанной таблицы перед выполнением последующих операций |
create_table() |
Создает новую таблицу с указанными семействами столбцов |
delete_table() |
Удаляет существующую таблицу |
put_row() |
Создает новую строку или обновляет существующую |
get_row() |
Получает одну строку по ее ключу |
delete_row() |
Удаляет всю строку или выбранные столбцы в строке |
scan_table() |
Сканирует диапазон строк и возвращает найденные записи |
batch_put_rows() |
Записывает несколько строк с использованием настраиваемого размера пакета и при необходимости параллельного выполнения |
batch_get_rows() |
Получает несколько строк в рамках одной пакетной операции |
batch_delete_rows() |
Удаляет несколько строк пакетами |
HBaseCLIHook
Хук HBaseCLIHook предоставляет методы, соответствующие операциям HBase Backup CLI.
| Метод | Описание |
|---|---|
create_full_backup() |
Создает полную резервную копию одной или нескольких таблиц HBase |
create_incremental_backup() |
Создает инкрементную резервную копию на основе предыдущего образа резервной копии |
restore_backup() |
Восстанавливает данные из существующей резервной копии |
get_backup_history() |
Получает историю завершенных операций резервного копирования |
describe_backup() |
Отображает подробную информацию об образе резервной копии |
create_backup_set() |
Создает или изменяет набор резервного копирования, используемый при выполнении операций резервного копирования |
list_backup_sets() |
Возвращает настроенные наборы резервного копирования |
execute_command() |
Выполняет произвольную команду HBase CLI и возвращает результат ее выполнения |
Операторы
Операторы реализуют функции, которые можно использовать в DAG. Каждый оператор инкапсулирует определенную операцию HBase, проверяет переданные параметры, инициализирует соответствующий хук, выполняет действие и при необходимости возвращает результат в контекст задачи Airflow.
Операторы управления таблицами
Операторы управления таблицами выполняют административные операции с таблицами HBase. Провайдер включает операторы управления таблицами, представленные ниже.
| Оператор | Описание |
|---|---|
HBaseCreateTableOperator |
Создает таблицу HBase и при необходимости проверяет, существует ли таблица до ее создания |
HBaseDeleteTableOperator |
Удаляет существующую таблицу и при необходимости проверяет ее существование перед выполнением операции |
Операторы для работы с данными
Операторы для работы с данными выполняют операции чтения и записи через HBase Thrift2 Server API. Они используют HBaseThriftHook, который автоматически выбирает либо отдельное соединение Thrift, либо соединение из пула в зависимости от настроенного подключения Airflow.
Провайдер включает следующие операторы для обработки данных.
| Оператор | Описание |
|---|---|
HBasePutOperator |
Вставляет или обновляет одну строку |
HBaseBatchPutOperator |
Записывает несколько строк с использованием настраиваемой пакетной обработки и при необходимости параллельного выполнения |
HBaseBatchGetOperator |
Получает несколько строк за одну операцию |
HBaseScanOperator |
Сканирует содержимое таблицы в заданном диапазоне строк |
Операторы резервного копирования и восстановления
Операторы резервного копирования предоставляют доступ к подсистеме резервного копирования HBase через HBaseCLIHook. Провайдер включает операторы для управления резервным копированием, представленные ниже.
| Оператор | Описание |
|---|---|
HBaseCreateBackupOperator |
Создает полные или инкрементные резервные копии HBase |
HBaseBackupHistoryOperator |
Получает историю завершенных операций резервного копирования |
HBaseRestoreOperator |
Восстанавливает таблицы из существующей резервной копии |
HBaseBackupSetOperator |
Создает наборы резервного копирования HBase и управляет ими |
Сенсоры
Сенсоры позволяют DAG ожидать, пока в HBase не будет выполнено определенное условие, прежде чем продолжить работу. Все параметры сенсора, идентифицирующие ресурсы HBase, такие как имена таблиц и ключи строк, определяются как поля шаблонов Airflow. Это позволяет сенсорам отслеживать динамически создаваемые ресурсы, имена которых определяются во время выполнения DAG.
Провайдер включает следующие сенсоры.
| Сенсор | Описание |
|---|---|
HBaseTableSensor |
Ожидает, пока указанная таблица HBase не станет доступной. Во время каждого цикла проверки сенсор вызывает метод |
HBaseRowSensor |
Ожидает появления определенной строки в таблице HBase. Во время каждого цикла проверки сенсор получает строку с помощью метода |
Вспомогательные классы
Провайдер определяет несколько классов-перечислений, предоставляющих строго типизированные значения конфигурации для операторов HBase. Использование этих перечислений вместо строковых литералов повышает читаемость определений DAG, проверяет значения параметров до начала выполнения и снижает вероятность ошибок конфигурации.
Провайдер определяет следующие вспомогательные классы.
| Тип | Описание |
|---|---|
BackupType |
Определяет, выполняется ли операция резервного копирования как полное или инкрементное резервное копирование |
BackupSetAction |
Определяет действие, выполняемое |
IfExistsAction |
Определяет поведение операторов создания таблиц, если целевая таблица уже существует |
IfNotExistsAction |
Определяет поведение операторов удаления таблиц, если целевая таблица не существует |
Использование провайдера в DAG
Чтобы использовать провайдер HBase, импортируйте необходимые операторы или сенсоры из пакета провайдера и настройте подключение HBase.
Примеры в этом разделе демонстрируют наиболее распространенные сценарии использования провайдера. Полные примеры DAG доступны в репозитории провайдера.
Базовый поток обработки данных
Типичный рабочий процесс HBase состоит из создания таблицы, записи данных и чтения сохраненных записей.
Используйте HBaseCreateTableOperator, чтобы создать таблицу HBase и определить ее семейства столбцов:
from airflow.providers.arenadata.hbase.operators.hbase import HBaseCreateTableOperator
create_table = HBaseCreateTableOperator(
task_id="create_table",
table_name="users", (1)
families={ (2)
"profile": {},
"contacts": {},
},
hbase_conn_id="hbase_default",
)
| 1 | Имя таблицы HBase. |
| 2 | Семейства столбцов, создаваемые вместе с таблицей. |
Используйте HBasePutOperator, чтобы вставить или обновить одну строку:
from airflow.providers.arenadata.hbase.operators.hbase import HBasePutOperator
insert_user = HBasePutOperator(
task_id="insert_user",
table_name="users",
row_key="user_001", (1)
data={ (2)
"profile:name": "Alice",
"profile:age": "30",
"contacts:email": "alice@example.com",
},
hbase_conn_id="hbase_default",
)
| 1 | Уникальный идентификатор строки. |
| 2 | Имена столбцов включают префикс семейства столбцов. |
Используйте HBaseScanOperator, чтобы получить строки из таблицы:
from airflow.providers.arenadata.hbase.operators.hbase import HBaseScanOperator
scan_users = HBaseScanOperator(
task_id="scan_users",
table_name="users",
row_start="user_000",
row_stop="user_999",
columns=[
"profile:name",
"contacts:email",
],
limit=100,
hbase_conn_id="hbase_default",
)
Эти операторы можно объединить в простой рабочий процесс загрузки данных.
Пакетная обработка
Для больших наборов данных вы можете использовать операторы пакетной обработки, чтобы уменьшить количество запросов Thrift и повысить пропускную способность.
HBaseBatchPutOperator записывает несколько строк с использованием настраиваемого размера пакета:
from airflow.providers.arenadata.hbase.operators.hbase import HBaseBatchPutOperator
batch_insert = HBaseBatchPutOperator(
task_id="batch_insert",
table_name="users",
rows=rows,
batch_size=200, (1)
max_workers=4, (2)
hbase_conn_id="hbase_default",
)
| 1 | Количество строк, записываемых за один запрос. |
| 2 | Количество параллельных рабочих потоков, используемых для пакетной обработки. |
При использовании нескольких рабочих потоков включите пул соединений в подключении HBase, чтобы повысить пропускную способность.
Используйте HBaseBatchGetOperator, чтобы получить несколько строк за одну операцию:
from airflow.providers.arenadata.hbase.operators.hbase import HBaseBatchGetOperator
batch_get = HBaseBatchGetOperator(
task_id="batch_get",
table_name="users",
row_keys=[
"user_001",
"user_002",
"user_003",
],
columns=[
"profile:name",
"contacts:email",
],
hbase_conn_id="hbase_default",
)
Рабочий процесс резервного копирования
Провайдер интегрируется с HBase Backup CLI для управления наборами резервного копирования и выполнения операций резервного копирования и восстановления.
Перед созданием резервных копий подготовьте директорию резервного копирования в HDFS:
$ hdfs dfs -mkdir -p /hbase/backup
$ hdfs dfs -chmod 777 /hbase/backup
Права доступа chmod 777 используются только в целях тестирования и не подходят для использования в производственной среде.
Включите поддержку резервного копирования в конфигурации HBase:
hbase.backup.enable=true
Набор резервного копирования объединяет одну или несколько таблиц в повторно используемый объект резервного копирования.
Создайте набор резервного копирования:
from airflow.providers.arenadata.hbase.operators.hbase import (
BackupSetAction,
HBaseBackupSetOperator,
)
create_backup_set = HBaseBackupSetOperator(
task_id="create_backup_set",
action=BackupSetAction.ADD,
backup_set_name="production_tables",
tables=[
"users",
"orders",
],
hbase_conn_id="hbase_default",
)
Используйте HBaseCreateBackupOperator, чтобы создать полную или инкрементную резервную копию.
Создайте резервную копию:
from airflow.providers.arenadata.hbase.operators.hbase import (
BackupType,
HBaseCreateBackupOperator,
)
create_backup = HBaseCreateBackupOperator(
task_id="create_backup",
backup_type=BackupType.FULL,
backup_path="hdfs:///hbase/backup", (1)
backup_set_name="production_tables",
workers=2,
hbase_conn_id="hbase_default",
)
| 1 | Расположение резервной копии в HDFS. |
Оператор передает сгенерированный идентификатор резервной копии в XCom для последующего использования при восстановлении.
Используйте HBaseBackupHistoryOperator, чтобы проверить завершенные операции резервного копирования:
from airflow.providers.arenadata.hbase.operators.hbase import HBaseBackupHistoryOperator
backup_history = HBaseBackupHistoryOperator(
task_id="backup_history",
backup_set_name="production_tables",
hbase_conn_id="hbase_default",
)
Восстановите данные, используя идентификатор резервной копии, возвращенный задачей резервного копирования:
from airflow.providers.arenadata.hbase.operators.hbase import HBaseRestoreOperator
restore_backup = HBaseRestoreOperator(
task_id="restore_backup",
backup_path="hdfs:///hbase/backup",
backup_id="{{ ti.xcom_pull(task_ids='create_backup') }}",
tables=["users"],
overwrite=True,
hbase_conn_id="hbase_default",
)
Мониторинг
Используйте HBaseTableSensor, чтобы ожидать появления таблицы:
from airflow.providers.arenadata.hbase.sensors.hbase import HBaseTableSensor
wait_for_table = HBaseTableSensor(
task_id="wait_for_table",
table_name="users",
timeout=300,
poke_interval=30,
hbase_conn_id="hbase_default",
)
Этот сенсор обычно используется, когда таблица создается другим DAG или внешним приложением.
Используйте HBaseRowSensor, чтобы ожидать появления определенной строки:
from airflow.providers.arenadata.hbase.sensors.hbase import HBaseRowSensor
wait_for_row = HBaseRowSensor(
task_id="wait_for_row",
table_name="users",
row_key="user_001",
timeout=600,
poke_interval=60,
hbase_conn_id="hbase_default",
)
Этот сенсор обычно используется для ожидания появления строк-маркеров, указывающих на завершение загрузки данных.