Airflow и Amazon EMR: Настройка AWS Provider и операторов для создания потоков заданий

Apache Airflow – мощный инструмент оркестрации рабочих процессов, который позволяет автоматизировать сложные процессы обработки данных. Amazon EMR (Elastic MapReduce) – облачный сервис, предоставляющий платформу для обработки больших данных с использованием фреймворков, таких как Spark, Hadoop и Hive. Интеграция Airflow и EMR позволяет автоматизировать запуск, мониторинг и завершение EMR кластеров, а также выполнение задач обработки данных.

В этой статье мы рассмотрим, как настроить AWS Provider в Airflow, какие операторы доступны для взаимодействия с EMR, и как создать DAG (Directed Acyclic Graph) для запуска задач на EMR кластере. Также будут рассмотрены расширенные сценарии и лучшие практики.

Подготовка к работе: Настройка AWS Provider в Airflow

Для взаимодействия Airflow с ресурсами AWS необходимо настроить AWS Provider. Провайдер предоставляет интерфейс для аутентификации и авторизации Airflow для доступа к сервисам AWS, таким как EMR.

Установка и настройка AWS Provider в Airflow

AWS Provider устанавливается с помощью pip:

pip install apache-airflow-providers-amazon

После установки необходимо настроить подключение к AWS. Это можно сделать несколькими способами, включая использование переменных окружения, IAM ролей или файла конфигурации AWS.

Аутентификация и авторизация Airflow для доступа к AWS ресурсам (IAM роли)

Рекомендуемый способ аутентификации – использование IAM ролей. IAM роль позволяет Airflow получать временные учетные данные для доступа к ресурсам AWS. Для этого необходимо создать IAM роль с необходимыми разрешениями для работы с EMR (например, AmazonEMRFullAccess) и назначить ее инстансу Airflow.

Пример политики IAM:

{
    "Version": "2012-10-17",
    "Statement": [
        {
            "Effect": "Allow",
            "Action": [
                "elasticmapreduce:*",
                "ec2:Describe*",
                "s3:*",
                "iam:PassRole"
            ],
            "Resource": "*"
        }
    ]
}

В Airflow необходимо настроить подключение AWS, указав aws_conn_id. Если Airflow запущен на инстансе EC2 с IAM ролью, Airflow автоматически использует IAM роль для аутентификации.

Операторы Airflow для взаимодействия с Amazon EMR

Airflow предоставляет набор операторов для взаимодействия с Amazon EMR. Эти операторы позволяют создавать, запускать, мониторить и завершать EMR кластеры, а также выполнять шаги (steps) обработки данных.

Обзор основных EMR операторов: EMRCreateJobFlowOperator, EMRStepSensor, EMRTerminateJobFlowOperator и др.

Основные операторы Airflow для работы с EMR:

  • EMRCreateJobFlowOperator: Создает EMR кластер.

  • EMRTerminateJobFlowOperator: Завершает EMR кластер.

  • EMRAddStepsOperator: Добавляет шаги (steps) для выполнения на EMR кластере.

  • EMRStepSensor: Ожидает завершения шага на EMR кластере.

  • EmrContainerOperator: Запускает контейнерные приложения на EMR в EKS.

Пример использования EMRCreateJobFlowOperator:

from airflow.providers.amazon.aws.operators.emr_create_job_flow import EmrCreateJobFlowOperator

JOB_FLOW_OVERRIDES = {
    "name": "emr-cluster-from-airflow",
    "releaseLabel": "emr-6.9.0",
    "applications": [{"Name": "Spark"}],
    "instances": {
        "instanceGroups": [
            {
                "name": "Master nodes",
                "market": "ON_DEMAND",
                "instanceRole": "MASTER",
                "instanceType": "m5.xlarge",
                "instanceCount": 1,
            },
            {
                "name": "Worker nodes",
                "market": "ON_DEMAND",
                "instanceRole": "CORE",
                "instanceType": "m5.xlarge",
                "instanceCount": 2,
            },
        ],
        "keepJobFlowAliveWhenNoSteps": True,
        "terminationProtected": False,
    },
    "jobFlowRole": "EMR_EC2_DefaultRole",
    "serviceRole": "EMR_DefaultRole",
}

create_emr_cluster = EmrCreateJobFlowOperator(
    task_id="create_emr_cluster",
    job_flow_overrides=JOB_FLOW_OVERRIDES,
    aws_conn_id="aws_default",
    region_name="us-east-1",
)
Реклама

Параметризация EMR операторов: передача параметров, использование XComs

Параметры для EMR операторов можно передавать непосредственно в коде DAG, а также использовать XComs для передачи данных между задачами. XComs позволяют передавать значения, такие как идентификатор кластера или результаты выполнения шагов, между задачами в DAG.

Пример использования XComs:

from airflow.providers.amazon.aws.operators.emr_create_job_flow import EmrCreateJobFlowOperator
from airflow.providers.amazon.aws.sensors.emr import EmrJobFlowSensor
from airflow.utils.dates import days_ago

from airflow import DAG

with DAG(dag_id='emr_with_xcoms', start_date=days_ago(2), schedule_interval=None, tags=['emr']) as dag:
    create_emr_cluster = EmrCreateJobFlowOperator(
        task_id='create_emr_cluster',
        job_flow_overrides=JOB_FLOW_OVERRIDES,
        aws_conn_id='aws_default',
        region_name='us-east-1'
    )

    check_cluster_status = EmrJobFlowSensor(
        task_id='check_cluster_status',
        job_flow_id= '{{ task_instance.xcom_pull("create_emr_cluster", key="return_value") }}',
        aws_conn_id='aws_default',
        region_name='us-east-1'
    )

    create_emr_cluster >> check_cluster_status

Создание DAG для запуска EMR кластера и выполнения заданий

Создание DAG для запуска EMR кластера и выполнения заданий включает несколько этапов: создание кластера, добавление шагов, мониторинг выполнения шагов и завершение кластера.

Пошаговое руководство: разработка DAG для запуска Spark/Hadoop job на EMR

  1. Определение параметров кластера: Определите параметры EMR кластера, такие как версия EMR, типы инстансов, количество узлов и список приложений (Spark, Hadoop, Hive и т.д.).

  2. Создание DAG: Создайте DAG в Airflow, определив зависимости между задачами.

  3. Создание кластера: Используйте EMRCreateJobFlowOperator для создания EMR кластера.

  4. Добавление шагов: Используйте EMRAddStepsOperator для добавления шагов обработки данных (например, запуск Spark job). Параметры шага можно передавать через JSON файл или непосредственно в коде DAG.

  5. Мониторинг шагов: Используйте EMRStepSensor для мониторинга выполнения шагов. Sensor ожидает завершения шага и переходит к следующей задаче.

  6. Завершение кластера: Используйте EMRTerminateJobFlowOperator для завершения EMR кластера.

Управление жизненным циклом EMR кластера: создание, запуск задач и завершение кластера

Airflow позволяет полностью автоматизировать жизненный цикл EMR кластера, начиная от создания и заканчивая завершением. Это позволяет оптимизировать затраты на ресурсы AWS и упростить управление процессами обработки данных.

Расширенные сценарии и лучшие практики

Примеры использования Airflow для ETL процессов на EMR

Airflow можно использовать для ETL (Extract, Transform, Load) процессов на EMR. Например, можно настроить DAG, который извлекает данные из S3, преобразует их с помощью Spark на EMR кластере, и загружает результат в Redshift или другую базу данных.

Мониторинг и обработка ошибок при работе с EMR и Airflow

Важно настроить мониторинг выполнения задач на EMR и обработку ошибок. Airflow предоставляет возможности для отслеживания состояния задач, отправки уведомлений об ошибках и повторного запуска задач.

Для мониторинга можно использовать:

  • Логи Airflow и EMR.

  • Метрики CloudWatch.

  • Интеграцию с системами мониторинга, такими как Prometheus и Grafana.

Обработка ошибок может включать:

  • Повторный запуск задач (retries).

  • Отправку уведомлений об ошибках (email, Slack).

  • Прерывание DAG в случае критических ошибок.

Заключение

Интеграция Airflow и Amazon EMR предоставляет мощный инструмент для оркестрации и автоматизации процессов обработки данных. Настройка AWS Provider, использование EMR операторов и создание DAG позволяют эффективно управлять EMR кластерами и выполнять сложные задачи обработки данных. Следуя лучшим практикам и настроив мониторинг и обработку ошибок, можно добиться высокой надежности и эффективности ETL процессов на EMR.


Добавить комментарий