Провайдер Ozone
Обзор
Провайдер Ozone позволяет DAG взаимодействовать с Apache Ozone. Он предоставляет хуки, операторы, сенсоры и операторы передачи данных для выполнения административных операций, управления файлами и директориями, а также передачи данных между Ozone и другими системами хранения.
Провайдер можно использовать для выполнения следующих операций в Ozone-хранилище:
-
создание и удаление томов и бакетов;
-
загрузка, скачивание, перемещение, копирование и удаление файлов;
-
управление директориями;
-
передача данных из HDFS в Ozone;
-
создание резервных копий на основе снепшотов;
-
мониторинг доступности файлов и директорий.
Провайдер соответствует стандартной архитектуре провайдеров в Airflow, в которой операторы реализуют задачи DAG, хуки обеспечивают низкоуровневое взаимодействие с внешними системами, сенсоры ожидают выполнения заданного условия, операторы передачи данных перемещают данные между системами хранения, а вспомогательные модули предоставляют функции, которые можно использовать повторно.
|
ВАЖНО
Провайдер Ozone не реализует S3-совместимый API и не использует Ozone S3 Gateway.
|
Требования
Для работы провайдера Ozone требуется следующее:
-
Airflow 2.10.3 или более поздняя версия;
-
подключение Ozone, настроенное в Airflow;
-
клиент Ozone, установленный на хостах с воркерами Airflow;
-
сетевое подключение к Ozone.
Для работы с оператором HdfsToOzoneOperator необходимо выполнить следующие дополнительные требования:
-
на хостах с воркерами Airflow должен быть установлен клиент HDFS;
-
Hadoop DistCp должен быть доступен;
-
настроено сетевое подключение к HDFS.
Если у вас уже есть кластер ADH, необходимые компоненты клиента Ozone и клиента HDFS можно предоставить кластеру ADO через совместное использование хостов.
На хостах Airflow требуются только компоненты клиентов. Компоненты Ozone Manager, Storage Container Manager и DataNode не требуются.
Архитектура
Провайдер Ozone позволяет DAG взаимодействовать с Ozone с помощью нативных инструментов командной строки Ozone. Он предоставляет хуки (hook), операторы (operator), сенсоры (sensor) и операторы передачи данных (transfer operator) для управления хранилищем Ozone, мониторинга объектов файловой системы и передачи данных между HDFS и Ozone.
Провайдер состоит из нескольких модулей Python:
-
Hooks
Хуки реализуют уровень взаимодействия между Airflow и Ozone. Они подготавливают среду выполнения, применяют параметры подключения и выполняют нативные команды Ozone.
Провайдер включает следующие реализации хуков:
-
OzoneAdminHook— выполняет административные операции, такие как создание томов и бакетов, настройка квот и управление снепшотами. -
OzoneFsHook— выполняет операции файловой системы, включая загрузку, скачивание, копирование, перемещение, вывод списка и удаление объектов.
-
-
Operators
Операторы реализуют исполняемые задачи Airflow. Каждый оператор проверяет переданные параметры, инициализирует соответствующий хук, выполняет запрошенную операцию и возвращает результат выполнения в Airflow.
-
Sensors
Сенсоры периодически проверяют состояние ресурсов Ozone до тех пор, пока не будет выполнено указанное условие или не истечет настроенное время ожидания. Типичный сценарий использования — ожидание появления файла, созданного другим рабочим процессом, перед запуском последующей обработки.
-
Transfers
Операторы передачи данных обеспечивают интеграцию Ozone с внешними системами хранения. Текущий провайдер включает операторы для копирования данных из HDFS в Ozone с помощью Hadoop DistCp и создания резервных копий Ozone на основе снепшотов.
Во время выполнения DAG эти модули работают совместно. Airflow планирует выполнение оператора, который инициализирует соответствующий хук, используя настроенное подключение Airflow. Хук подготавливает среду выполнения, вызывает необходимую команду Ozone CLI и возвращает результат выполнения оператору. Затем оператор передает статус выполнения задачи обратно планировщику Airflow.
Конфигурация
Провайдер использует тип подключения Airflow Ozone для хранения параметров, необходимых для взаимодействия с кластером Ozone.
Подключение можно настроить в пользовательском интерфейсе Airflow.
Стандартная конфигурация
Провайдер использует стандартные поля подключения Airflow, приведенные ниже. Специфичные для провайдера параметры конфигурации хранятся в поле Extra в формате JSON.
| Поле | Описание |
|---|---|
Connection type |
Тип подключения, предоставляемый провайдером. Должен иметь значение |
Host |
Имя хоста или IP-адрес Ozone Manager |
Port |
Порт Ozone Manager |
Extra |
Конфигурация среды выполнения провайдера |
Пример базового подключения:
-
Connection Type:
ozone -
Host:
ozone-manager.example.com -
Port:
9862 -
Extra:
{
"ozone_conf_dir": "/etc/ozone/conf"
}
Дополнительные параметры
Провайдер использует дополнительные параметры среды выполнения из поля Extra.
| Параметр | Описание | Значение по умолчанию |
|---|---|---|
ozone_conf_dir |
Директория, содержащая файлы конфигурации клиентов Ozone/HDFS (ozone-site.xml и core-site.xml). Провайдер использует значение для установки переменных |
— |
ozone_security_enabled |
Включает режим безопасности Ozone в среде выполнения провайдера. Установите значение |
false |
ozone_om_https_port |
HTTPS-порт, используемый Ozone OM в кластере ADH (например, |
— |
hadoop_security_authentication |
Режим аутентификации. Установите значение |
simple |
max_content_size_bytes |
Максимальный размер загружаемого или скачиваемого содержимого |
1073741824 (1 GiB) |
SSL/TLS
Провайдер поддерживает шифрование при взаимодействии с кластером Ozone. Параметры SSL/TLS указываются в поле Extra подключения.
Следующий пример настраивает SSL-подключение:
{
"ozone_security_enabled": "true",
"ozone_om_https_port": "9879",
"ozone_ssl_keystore_location": "/etc/ssl/ozone-keystore.jks",
"ozone_ssl_keystore_password": "secret://vault/ozone/keystore_password",
"ozone_ssl_keystore_type": "JKS",
"ozone_ssl_truststore_location": "/etc/ssl/ozone-truststore.jks",
"ozone_ssl_truststore_password": "secret://vault/ozone/truststore_password",
"ozone_ssl_truststore_type": "JKS",
"ozone_conf_dir": "/etc/ozone/conf"
}
Доступные параметры SSL приведены ниже.
| Параметр | Описание | Значение по умолчанию |
|---|---|---|
ozone_ssl_keystore_location |
Путь к SSL-хранилищу ключей |
— |
ozone_ssl_keystore_password |
Пароль хранилища ключей |
— |
ozone_ssl_keystore_type |
Формат хранилища (обычно |
JKS |
ozone_ssl_truststore_location |
Путь к SSL-хранилищу доверенных сертификатов |
— |
ozone_ssl_truststore_password |
Пароль хранилища доверенных сертификатов |
— |
ozone_ssl_truststore_type |
Формат хранилища доверенных сертификатов |
JKS |
Kerberos
Провайдер поддерживает аутентификацию Kerberos для защищенных кластеров с Ozone.
Настройте параметры Kerberos в поле Extra подключения:
{
"hadoop_security_authentication": "kerberos",
"kerberos_principal": "airflow@EXAMPLE.COM",
"kerberos_keytab": "/etc/security/keytabs/airflow.keytab",
"krb5_conf": "/etc/krb5.conf",
"ozone_conf_dir": "/etc/ozone/conf"
}
Поддерживаются следующие параметры Kerberos.
| Параметр | Описание | Значение по умолчанию |
|---|---|---|
kerberos_principal |
Принципал Kerberos, используемый задачами Airflow |
— |
kerberos_keytab |
Путь к файлу Kerberos keytab |
— |
kerberos_password |
Пароль, используемый вместо keytab |
— |
krb5_conf |
Путь к файлу krb5.conf |
Системное значение по умолчанию |
Если включена аутентификация Kerberos, необходимо указать kerberos_keytab или kerberos_password.
{
"ozone_security_enabled": "true",
"ozone_om_https_port": "9879",
"ozone_ssl_keystore_location": "/opt/airflow/ssl/ozone-keystore.jks",
"ozone_ssl_keystore_password": "secret://vault/ozone/keystore_password",
"ozone_ssl_keystore_type": "JKS",
"ozone_ssl_truststore_location": "/opt/airflow/ssl/ozone-truststore.jks",
"ozone_ssl_truststore_password": "secret://vault/ozone/truststore_password",
"ozone_ssl_truststore_type": "JKS",
"hadoop_security_authentication": "kerberos",
"kerberos_principal": "testuser@EXAMPLE.COM",
"kerberos_keytab": "/opt/airflow/keytabs/testuser.keytab",
"krb5_conf": "/opt/airflow/kerberos-config/krb5.conf",
"ozone_conf_dir": "/opt/airflow/ozone-conf"
}
{
"ozone_security_enabled": "true",
"ozone_om_https_port": "9879",
"ozone_ssl_keystore_location": "/opt/airflow/ssl/ozone-keystore.jks",
"ozone_ssl_keystore_password": "secret://vault/ozone/keystore_password",
"ozone_ssl_keystore_type": "JKS",
"ozone_ssl_truststore_location": "/opt/airflow/ssl/ozone-truststore.jks",
"ozone_ssl_truststore_password": "secret://vault/ozone/truststore_password",
"ozone_ssl_truststore_type": "JKS",
"hadoop_security_authentication": "kerberos",
"kerberos_principal": "testuser@EXAMPLE.COM",
"kerberos_password": "secret://vault/ozone/kerberos_password",
"krb5_conf": "/opt/airflow/kerberos-config/krb5.conf",
"ozone_conf_dir": "/opt/airflow/ozone-conf"
}
Конфигурация передачи данных HDFS
Оператор HdfsToOzoneOperator использует подключение HDFS в дополнение к подключению Ozone. Параметры безопасности, необходимые для Hadoop DistCp, считываются из подключения, указанного в hdfs_conn_id.
Если для HDFS включена аутентификация Kerberos, настройте подключение HDFS, указав соответствующий принципал и либо keytab, либо пароль.
Ограничения на загрузку и скачивание
Провайдер ограничивает максимальный размер данных для загрузки с помощью параметра подключения max_content_size_bytes.
Операторы загрузки и скачивания также поддерживают параметр max_content_size_bytes, который переопределяет значение, заданное в подключении, для отдельной задачи.
Хуки
Хуки реализуют уровень взаимодействия провайдера. Поскольку провайдер Ozone использует модель выполнения на основе CLI, хуки подготавливают и настраивают среду выполнения команд, вызывают нативные команды CLI и возвращают результаты выполнения операторам и сенсорам.
Провайдер включает следующие хуки.
| Хук | Описание |
|---|---|
OzoneCliHook |
Базовый хук, отвечающий за разбор подключений Airflow, подготовку среды выполнения, применение параметров SSL/TLS и Kerberos, выполнение нативных команд Ozone CLI с повторными попытками и тайм-аутами, а также проверку подключений Airflow |
OzoneFsHook |
Расширяет |
OzoneAdminHook |
Расширяет |
OzoneAdminExtraHook |
Расширяет |
Операторы
Операторы реализуют функциональность внутри задач Airflow, которые позволяют выполнять операции с Ozone. Каждый оператор проверяет свои параметры, инициализирует соответствующий хук, выполняет требуемую команду CLI и возвращает результат выполнения в Airflow.
В зависимости от выполняемой операции, используется либо хук OzoneAdminHook, либо OzoneFsHook. При необходимости операторы передачи данных дополнительно вызывают утилиты Hadoop DistCp.
Административные операторы
Административные операторы управляют томами, бакетами и квотами Ozone.
| Оператор | Описание |
|---|---|
OzoneCreateVolumeOperator |
Создает том Ozone |
OzoneCreateBucketOperator |
Создает бакет в существующем томе |
OzoneSetQuotaOperator |
Применяет квоту хранилища к тому или бакету |
OzoneDeleteVolumeOperator |
Удаляет существующий том |
OzoneDeleteBucketOperator |
Удаляет существующий бакет |
Операторы файловой системы
Операторы файловой системы выполняют команды ozone fs для управления директориями и объектами, хранящимися в Ozone.
| Оператор | Описание |
|---|---|
OzoneCreatePathOperator |
Создает директорию в файловой системе Ozone |
OzoneUploadContentOperator |
Загружает встроенное текстовое содержимое как новый объект |
OzoneUploadFileOperator |
Загружает локальный файл в Ozone |
OzoneDownloadFileOperator |
Скачивает объект из Ozone в локальную файловую систему |
OzoneDeleteKeyOperator |
Удаляет один или несколько объектов |
OzoneDeletePathOperator |
Удаляет директорию |
OzoneListOperator |
Выводит список объектов в директории |
OzoneMoveOperator |
Перемещает или переименовывает объект |
OzoneCopyOperator |
Копирует объект внутри Ozone |
OzonePathExistsOperator |
Проверяет существование файла или директории |
Большинство операторов файловой системы поддерживают настраиваемое поведение, если целевой объект уже существует. В зависимости от выбранной политики оператор может завершить работу с ошибкой, проигнорировать существующий объект или перезаписать загружаемые файлы.
Операторы загрузки и скачивания также поддерживают настраиваемые ограничения размера с помощью параметра max_content_size_bytes. Если этот параметр явно не указан, провайдер использует значение, заданное в подключении Airflow.
Операторы передачи данных
Операторы передачи данных обеспечивают интеграцию Ozone с внешними системами хранения и средствами резервного копирования.
| Оператор | Описание |
|---|---|
HdfsToOzoneOperator |
Передает данные из HDFS в Ozone с помощью Hadoop DistCp |
OzoneBackupOperator |
Создает снимок бакета Ozone для резервного копирования |
HdfsToOzoneOperator считывает параметры безопасности HDFS из подключения, указанного в hdfs_conn_id, и применяет их только к среде выполнения DistCp.
Сенсоры
Сенсоры позволяют DAG ожидать, пока ресурс не станет доступен, прежде чем продолжить выполнение рабочего процесса.
В настоящее время провайдер включает следующий сенсор.
| Сенсор | Описание |
|---|---|
OzoneKeySensor |
Ожидает, пока указанный объект или каталог не станет доступен в файловой системе Ozone |
OzoneKeySensor использует OzoneFsHook для периодической проверки существования указанного пути.
Сенсор поддерживает шаблонизируемые пути, позволяя DAG отслеживать динамически создаваемые пути. Во время каждого цикла опроса сенсор выполняет быструю проверку существования через CLI, используя ту же конфигурацию среды выполнения и параметры безопасности, что и остальные компоненты провайдера.
В отличие от стандартного тайм-аута сенсора Airflow, параметр timeout в OzoneKeySensor определяет тайм-аут одного выполнения CLI. Общее поведение сенсора, включая частоту опроса и максимальное время ожидания, управляется стандартными параметрами сенсоров Airflow, такими как poke_interval, а также параметрами, унаследованными от BaseSensorOperator.
Использование провайдера в DAG
Чтобы использовать провайдер Ozone, импортируйте необходимые операторы или сенсоры из пакета провайдера и настройте подключение Ozone.
Провайдер включает несколько примеров DAG, демонстрирующих типичные административные сценарии и сценарии обработки данных. Эти примеры можно выполнять без изменения исходного кода DAG. Параметры среды выполнения можно переопределить в диалоговом окне Airflow Trigger DAG или передать через dag_run.conf.
Базовый сценарий работы с хранилищем
DAG example_ozone_usage демонстрирует базовый сценарий работы с Ozone. Он создает том и бакет, загружает файл в Ozone и проверяет существование загруженного объекта.
Пример создания тома и бакета:
from airflow.providers.arenadata.ozone.operators.ozone import (
OzoneCreateBucketOperator,
OzoneCreateVolumeOperator,
)
create_volume = OzoneCreateVolumeOperator(
task_id="create_volume",
volume_name="vol1", (1)
quota="10GB", (2)
ozone_conn_id="ozone_admin_default", (3)
)
create_bucket = OzoneCreateBucketOperator(
task_id="create_bucket",
volume_name="vol1", (4)
bucket_name="bucket-native", (5)
quota="10GB",
ozone_conn_id="ozone_admin_default",
)
| 1 | Имя создаваемого тома. |
| 2 | Квота хранилища тома. |
| 3 | Подключение Airflow, используемое для взаимодействия с Ozone. |
| 4 | Целевой том. |
| 5 | Бакет, создаваемый внутри указанного тома. |
Пример загрузки файла в Ozone:
from airflow.providers.arenadata.ozone.operators.ozone import OzoneUploadContentOperator
upload_file = OzoneUploadContentOperator(
task_id="upload_file",
content="Hello from Airflow",
remote_path="ofs://om-service/vol1/bucket-native/data_dir/file.txt", (1)
ozone_conn_id="ozone_admin_default",
if_exists="overwrite", (2)
)
| 1 | Целевой путь в файловой системе Ozone. |
| 2 | Заменяет существующий объект, позволяя многократно выполнять пример DAG. |
Пример проверки того, что загруженный файл существует:
from airflow.providers.arenadata.ozone.sensors.ozone import OzoneKeySensor
wait_for_file = OzoneKeySensor(
task_id="wait_for_file",
path="ofs://om-service/vol1/bucket-native/data_dir/file.txt", (1)
ozone_conn_id="ozone_admin_default",
mode="reschedule", (2)
)
| 1 | Путь к отслеживаемому объекту. |
| 2 | Освобождает слот воркера между попытками. |
Управление жизненным циклом данных
DAG example_ozone_data_lifecycle демонстрирует полный сценарий управления жизненным циклом данных, включая обнаружение файлов, обработку, архивирование, резервное копирование и очистку.
Пример вывода списка файлов, хранящихся в бакете:
from airflow.providers.arenadata.ozone.operators.ozone import OzoneListOperator
list_files = OzoneListOperator(
task_id="list_files_in_landing_zone",
path="ofs://om-service/landing/raw/*", (1)
ozone_conn_id="ozone_admin_default",
)
| 1 | Выводит список всех объектов, расположенных по пути. Результат сохраняется в XCom для последующих задач. |
Пример архивации обработанных файлов:
from airflow.providers.arenadata.ozone.operators.ozone import OzoneMoveOperator
archive_files = OzoneMoveOperator(
task_id="archive_landing_files",
source_path="ofs://om-service/landing/raw/*", (1)
dest_path="ofs://om-service/archive/processed/ds={{ ds }}", (2)
ozone_conn_id="ozone_admin_default",
)
| 1 | Перемещает все обработанные файлы из директории. |
| 2 | Сохраняет архивные данные в директории, партиционированной по дате. |
Пример удаления исходных файлов:
from airflow.providers.arenadata.ozone.operators.ozone import OzoneDeleteKeyOperator
cleanup = OzoneDeleteKeyOperator(
task_id="cleanup_landing_zone",
path="ofs://om-service/landing/raw/*",
ozone_conn_id="ozone_admin_default",
)
Этот сценарий можно использовать как шаблон для автоматизированных пайплайнов архивирования и хранения данных.
Административный сценарий
DAG example_ozone_multi_tenant_management автоматизирует подготовку ресурсов хранилища для новых проектов или тенантов.
Пример создания тома проекта и настройки квоты хранилища:
from airflow.providers.arenadata.ozone.operators.ozone import (
OzoneCreateVolumeOperator,
OzoneSetQuotaOperator,
)
create_volume = OzoneCreateVolumeOperator(
task_id="create_project_volume",
volume_name="project-alpha", (1)
ozone_conn_id="ozone_admin_default",
)
set_quota = OzoneSetQuotaOperator(
task_id="set_project_quota",
volume="project-alpha",
quota="10GB", (2)
ozone_conn_id="ozone_admin_default",
)
| 1 | Создает изолированный том хранилища для проекта. |
| 2 | Ограничивает максимальный объем хранилища, доступный для тома. |
Пример создания бакета проекта:
create_landing_bucket = OzoneCreateBucketOperator(
task_id="create_landing_bucket",
volume_name="project-alpha",
bucket_name="landing", (1)
quota="1GB",
ozone_conn_id="ozone_admin_default",
)
| 1 | Создает бакет для данных проекта внутри подготовленного тома. |
Передача данных
DAG example_ozone_copy_from_hdfs демонстрирует перенос данных из HDFS в Ozone с использованием Hadoop DistCp.
Сценарий выполняет следующие операции:
-
Ожидает появления файла в Ozone.
-
Устанавливает квоту для целевого бакета.
-
Копирует данные из HDFS в Ozone.
В демонстрационных целях DAG автоматически создает файл после небольшой задержки. В производственной среде эту вспомогательную ветку можно удалить, а файл может создаваться сторонним процессом загрузки данных.
Пример копирования данных из HDFS в Ozone:
from airflow.providers.arenadata.ozone.transfers.hdfs_to_ozone import HdfsToOzoneOperator
copy_hdfs_to_ozone = HdfsToOzoneOperator(
task_id="copy_hdfs_to_ozone",
source_path="hdfs:///user/data/legacy/", (1)
dest_path="ofs://om-service/vol1/bucket1/migrated/", (2)
hdfs_conn_id="hdfs_admin_default", (3)
)
| 1 | Исходная директория в HDFS. |
| 2 | Целевая директория в файловой системе Ozone. |
| 3 | Подключение Airflow, используемое для доступа к HDFS. |
Мониторинг
Используйте OzoneKeySensor, чтобы ожидать появления объекта в Ozone.
Этот сенсор обычно используется для синхронизации независимых рабочих процессов, обнаружения появления новых наборов данных или ожидания завершения записи данных другими сервисами.
from airflow.providers.arenadata.ozone.sensors.ozone import OzoneKeySensor
wait_for_trigger = OzoneKeySensor(
task_id="wait_for_trigger_file",
path="ofs://om-service/vol1/bucket1/trigger.lck", (1)
ozone_conn_id="ozone_admin_default",
mode="reschedule",
)
| 1 | Объект, наличие которого запускает выполнение последующих задач. |
Создание резервной копии
Пример создания резервной копии на основе снепшота:
from airflow.providers.arenadata.ozone.transfers.ozone_backup import OzoneBackupOperator
backup_archive = OzoneBackupOperator(
task_id="backup_archive_via_snapshot",
volume="archive", (1)
bucket="processed", (2)
snapshot_name="snap-{{ ds_nodash }}", (3)
ozone_conn_id="ozone_admin_default",
)
| 1 | Том, содержащий архивные данные. |
| 2 | Бакет, для которого создается снепшот. |
| 3 | Имя снепшота, сформированное из указанного префикса и даты выполнения DAG. |