Провайдер 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

Тип подключения, предоставляемый провайдером. Должен иметь значение Ozone

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_CONF_DIR и HADOOP_CONF_DIR

 — 

ozone_security_enabled

Включает режим безопасности Ozone в среде выполнения провайдера. Установите значение true для сценариев с SSL и SSL + Kerberos

false

ozone_om_https_port

HTTPS-порт, используемый Ozone OM в кластере ADH (например, 9879 в локальных средах). Должен соответствовать фактической конфигурации SSL Ozone

 — 

hadoop_security_authentication

Режим аутентификации. Установите значение kerberos, чтобы включить аутентификацию Kerberos

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)

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.

Пример конфигурации SSL/TLS и Kerberos с использованием keytab
{
  "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"
}
Пример конфигурации SSL/TLS и Kerberos с использованием пароля
{
  "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

Расширяет OzoneCliHook и предоставляет операции файловой системы через ozone fs, включая создание директорий, загрузку и скачивание файлов, копирование, перемещение, удаление, вывод списка и проверку существования объектов файловой системы

OzoneAdminHook

Расширяет OzoneCliHook и предоставляет административные операции через ozone sh, включая управление томами и бакетами, настройку квот, получение метаданных и вывод списка ресурсов

OzoneAdminExtraHook

Расширяет OzoneAdminHook, добавляя расширенные административные операции, такие как управление снепшотами, управление тенантами, настройка репликации бакетов, управление ссылками на бакеты и получение отчетов Storage Container Manager (SCM)

Операторы

Операторы реализуют функциональность внутри задач 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.

Сценарий выполняет следующие операции:

  1. Ожидает появления файла в Ozone.

  2. Устанавливает квоту для целевого бакета.

  3. Копирует данные из 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.
Нашли ошибку? Выделите текст и нажмите Ctrl+Enter чтобы сообщить о ней