Saltar al contenido
Introducción a Apache Airflow

Arquitectura y componentes de Airflow.

Introducción

Apache Airflow es una plataforma de flujo de trabajo de código abierto que permite a los desarrolladores crear, programar y monitorear flujos de trabajo complejos de manera eficiente y escalable.

La arquitectura de Airflow se basa en un enfoque de software como servicio (SaaS) que consiste en la ejecución de flujos de trabajo en una serie de operaciones ordenadas y automatizadas. Airflow se compone de varios componentes fundamentales:

  • Metadatos de almacenamiento: almacenamiento de información relacionada con los flujos de trabajo, tareas y ejecuciones.
  • Árbol de dependencias: estructura en la que se basan los flujos de trabajo programados por los desarrolladores.
  • Programador: componente responsable de programar la ejecución de los flujos de trabajo en un momento específico.
  • Motor de ejecución: motor que ejecuta los flujos de trabajo, ya sea en paralelo o secuencialmente.
  • Interfaz de usuario: interfaz web proporcionada por Airflow para la administración y monitoreo de flujos de trabajo.

Gracias a su arquitectura modular y a estos componentes, Airflow se ha convertido en una herramienta muy popular entre los desarrolladores para la programación de flujos de trabajo complejos y escalables.

Resumen

Una estructura organizada y detallada de los tres tipos de componentes en la arquitectura de Apache Airflow son:

1. Componentes principales:

Estos son los componentes esenciales que definen la lógica central de Apache Airflow.

  • Webserver: Interfaz de usuario web de Airflow que permite administrar y monitorear el flujo de trabajo de las tareas.
  • Scheduler: Componente central que planifica y ejecuta las tareas en función de las dependencias y horarios definidos en los DAG (directed acyclic graph).
  • Executor: Interfaz que conecta la planificación de tareas del scheduler con la ejecución real de tareas en el backend.
  • Database: Base de datos utilizada para almacenar metadatos, como la definición de DAG y las configuraciones de flujo de trabajo.

2. Componentes de backend:

Estos componentes proporcionan la infraestructura necesaria para ejecutar las tareas y determinan cómo se realizará la ejecución.

  • Workers: Procesos que ejecutan las tareas en un servidor específico.
  • Celery: Backend recomendado para Airflow, una plataforma de canalización de tareas distribuida basada en mensajes.
  • Kubernetes: Backend opcional que permite la orquestación de contenedores a través de clústeres.

3. Componentes de extensión:

Estos componentes adicionales permiten la personalización y extensión de Airflow.

  • Hooks: Generan notificaciones o realizan acciones en servicios externos.
  • Operators: Definen cómo se lleva a cabo una tarea específica. Ejemplos incluyen BashOperator, PythonOperator, etc.
  • Plugins: Agregan nuevas características a Airflow, como nuevos hooks, operadores y métodos de extensión personalizados.

Esta estructura organizada proporciona una visión clara de cómo están distribuidos y funcionan los componentes dentro de Apache Airflow, facilitando la comprensión de su arquitectura y su capacidad para manejar flujos de trabajo complejos de manera eficiente y escalable.

Aplicación teórica

Un ejemplo práctico de la arquitectura y principales componentes de Apache Airflow:

Supongamos que tenemos una empresa de e-commerce que necesita procesar grandes cantidades de datos cada día. Queremos utilizar Airflow para automatizar este proceso de ETL (extracción, transformación y carga) de datos.

La arquitectura de Airflow consta de tres componentes principales:

  1. La base de datos: Un sistema de gestión de base de datos (como PostgreSQL o MySQL) que mantiene una tabla para almacenar los metadatos de las tareas y DAGs en Airflow.
  2. El scheduler: Una tarea de Airflow que se encarga de planificar y ejecutar las tareas de acuerdo con la definición de los DAGs.
  3. Los workers: Los procesos de trabajo que ejecutan las tareas específicas de nuestro DAG.

Ahora, podemos definir un DAG para nuestro proyecto de e-commerce. Este DAG tendrá tres tareas:

  1. Extraer datos: Una tarea que lee datos de fuentes externas, como puede ser un API o una base de datos, y los carga en memoria.
  2. Transformar datos: Una tarea que procesa y transforma los datos extraídos para que puedan ser cargados correctamente.
  3. Cargar datos: Una tarea que carga los datos transformados en una base de datos de destino.

Podemos definir estas tareas en Python utilizando la colección de operadores proporcionados por Airflow para interactuar con diferentes sistemas informáticos. Por ejemplo, para la primera tarea, podemos utilizar el operador PythonOperator junto con llamadas a una API externa:


from airflow import DAG
from airflow.operators.python_operator import PythonOperator
from datetime import datetime

dag = DAG(
    'ecommerce_etl',
    description='Procesamiento de datos diarios para una tienda en línea',
    schedule_interval='@daily',
    start_date=datetime(2022, 2, 1),
)

def extract_data():
    # Código para leer datos de una API externa
    return data

extract_task = PythonOperator(
    task_id='extraer_datos',
    python_callable=extract_data,
    dag=dag,
)
    

Finalmente, debemos definir cómo se relacionan estas tareas. Podemos conectar nuestras tareas definidas a través de la definición del DAG, estableciendo las dependencias entre ellas. Por ejemplo:


extract_task >> transform_task >> load_task
    

De esta forma, Airflow planificará y ejecutará automáticamente las tareas de nuestro DAG a través de su scheduler y los workers adecuados, asegurando que todos los datos necesarios sean procesados y cargados de forma puntual y eficiente.

Aplicación práctica

En Apache Airflow, se puede identificar tres componentes principales:

  1. Webserver: Es el componente que permite la visualización y monitoreo del estado de las tareas. Aquí se pueden visualizar las tareas activas e históricas, su estado, códigos de salida, etc.
  2. Scheduler: El programador se encarga de orquestar las tareas según su planificación definida en el DAG (Directed Acyclic Graph). Se ejecuta en segundo plano, identificando el momento en que cada tarea debe ejecutarse y actualiza la base de datos.
  3. Base de Datos: Almacena información necesaria para la ejecución de las tareas y su estado. Airflow es compatible con varias bases de datos, incluyendo SQLite, PostgreSQL y MySQL.

Aquí te muestro un ejemplo práctico de la arquitectura de Airflow en Python:


from airflow import DAG
from datetime import datetime, timedelta
from airflow.operators.bash_operator import BashOperator
from airflow.operators.python_operator import PythonOperator

default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'start_date': datetime(2021, 1, 1),
    'email_on_failure': False,
    'email_on_retry': False,
    'retries': 1,
    'retry_delay': timedelta(minutes=5)
}

dag = DAG(
    'simple_dag_example',
    default_args=default_args,
    description='Un ejemplo simple de DAG',
    schedule_interval=timedelta(days=1)
)

# Definir tareas
t1 = BashOperator(
    task_id='print_hello',
    bash_command='echo "Hola mundo"',
    dag=dag
)

def print_world():
    print('Mundo')

t2 = PythonOperator(
    task_id='print_world',
    python_callable=print_world,
    dag=dag
)

# Definir dependencias de tareas
t2.set_upstream(t1)
    

Este ejemplo define un DAG (grafos acíclicos dirigidos) con dos tareas: print_hello y print_world. print_hello es una tarea Bash que imprimirá "¡Hola, mundo!" cuando se ejecute. print_world es una tarea Python que imprimirá "mundo" en la consola. El set_upstream dice que print_world se ejecutará después de que print_hello haya completado su ejecución. Esto se define para que el DAG se ejecute de manera ordenada y dependiendo del resultado de cada tarea.