Saltar al contenido
Introducción a Apache Airflow

Creación de tareas personalizadas y operadores en Airflow.

Introducción

Apache Airflow es una plataforma robusta de flujo de trabajo de código abierto que se distingue por su flexibilidad y capacidad para adaptarse a diversas necesidades organizativas. Aquí te explico con más detalle sobre la creación de tareas personalizadas y operadores en Airflow:

Creación de Tareas Personalizadas y Operadores en Apache Airflow

  1. Clase BaseOperator:

    • Descripción: En Airflow, las tareas se definen utilizando clases que heredan de la clase BaseOperator.
    • Funcionalidad: BaseOperator proporciona la estructura básica para definir cómo se ejecutará una tarea dentro de un DAG.
    • Personalización: Al heredar de BaseOperator, los desarrolladores pueden definir el comportamiento específico de una tarea, como acciones a ejecutar, manejo de errores y condiciones de éxito.
    • Flexibilidad: Esta estructura permite implementar tareas que pueden ejecutar código Python, comandos Bash, consultas SQL u operaciones en Java, según las necesidades del flujo de trabajo.
  2. Soporte Multilenguaje:

    • Python: Es el lenguaje más comúnmente utilizado para definir operadores personalizados en Airflow debido a su integración profunda y facilidad de uso.
    • Otros lenguajes: Airflow también soporta la creación de operadores en Java, SQL y Bash, lo que amplía las opciones para integrar sistemas y procesos heterogéneos dentro de los flujos de trabajo.
  3. Integración de Nuevos Operadores:

    • Personalización: Los equipos pueden agregar nuevos operadores al proyecto de Airflow para extender la funcionalidad básica.
    • Ejemplos: Esto puede incluir operadores específicos para interactuar con APIs, servicios en la nube, sistemas de bases de datos particulares, entre otros.
    • Comunidad y Contribuciones: La comunidad de Airflow es activa en la creación y contribución de nuevos operadores, lo que enriquece continuamente el ecosistema de la plataforma.

Ventajas de Crear Tareas Personalizadas y Operadores

  • Adaptabilidad: Permite adaptar los flujos de trabajo a los requisitos específicos de cada organización.
  • Reutilización: Facilita la reutilización de código y la estandarización de procesos complejos.
  • Flexibilidad: Amplía las capacidades de Airflow más allá de los operadores preconstruidos, cubriendo casos de uso especializados y únicos.
  • Escalabilidad: Contribuye a la creación de flujos de trabajo escalables y robustos que manejan diferentes tipos de operaciones y sistemas.

En resumen, la capacidad de crear tareas personalizadas y operadores en Apache Airflow es fundamental para aprovechar al máximo la flexibilidad y potencia de la plataforma. Esto permite a las organizaciones diseñar flujos de trabajo eficientes y adaptados a sus necesidades específicas, asegurando una automatización efectiva y escalable de procesos empresariales complejos.

Resumen

Creación de Tareas Personalizadas en Apache Airflow

En Apache Airflow, las tareas personalizadas se crean mediante la definición de nuevos operadores. Un operador es una clase de Python que define las acciones a tomar para realizar una tarea en un DAG (Directed Acyclic Graph) de Airflow. Para crear un nuevo operador, es necesario definir una clase que herede de BaseOperator. Esta clase debe implementar el método execute, que es el encargado de llevar a cabo la tarea que se quiere ejecutar. Dentro de este método, se pueden realizar tareas como leer y escribir archivos, ejecutar comandos de shell, hacer consultas a bases de datos, entre otras.

Además de execute, existen otros métodos que se pueden sobrescribir para proporcionar mayor funcionalidad al operador. Por ejemplo, on_failure_callback, que se ejecuta cuando la tarea falla, o on_retry_callback, que se ejecuta cuando una tarea se intenta ejecutar de nuevo.

Los operadores predefinidos en Airflow incluyen tareas como BashOperator (para ejecutar comandos de shell), PythonOperator (para ejecutar código Python), PostgresOperator (para ejecutar consultas en una base de datos PostgreSQL), entre otros. Si estas tareas no cubren todas las necesidades de un proyecto, se puede crear un operador personalizado para una tarea específica.

Para utilizar un operador personalizado en un DAG, se debe importar en el archivo que define el DAG y luego crear una instancia de la clase del operador, pasando los argumentos necesarios. Por ejemplo:


from airflow import DAG
from my_operators import MyOperator

dag = DAG('my_dag')

my_task = MyOperator(
    task_id='my_task',
    param1='value1',
    param2='value2',
    dag=dag
)

En este ejemplo, my_task es una instancia de la clase MyOperator, que se ha creado con los parámetros param1 y param2. Cuando se ejecute este DAG, se llamará al método execute de MyOperator, que debe realizar alguna tarea específica.

Aplicación teórica

Flujo de Trabajo en Airflow para Preparación de Datos

Supongamos que tenemos un flujo de trabajo en Airflow que debe realizar una tarea de preparación de datos para un modelo de aprendizaje automático. En este flujo de trabajo se deben descargar datos de una fuente externa, realizar varias transformaciones en los datos y cargarlos en un sistema de almacenamiento para su posterior uso.

Para esta tarea, podríamos crear una tarea personalizada en Airflow llamada "descarga_datos" que use el operador "PythonOperator" para descargar los datos de una fuente externa. Podríamos crear otro operador personalizado llamado "transformacion_datos" que realice varias transformaciones en los datos descargados y cargue los datos transformados en un sistema de almacenamiento.

Aquí te muestro un ejemplo de cómo podría ser la definición de estas dos tareas personalizadas y operadores en Python:


from airflow import DAG
from airflow.operators.python_operator import PythonOperator
from mi_libreria import descargar_datos, transformacion_datos

dag = DAG(dag_id='preparacion_datos', schedule_interval='0 0 * * *')

descarga_datos_task = PythonOperator(
    task_id='descarga_datos',
    python_callable=descargar_datos,
    dag=dag
)

transformacion_datos_task = PythonOperator(
    task_id='transformacion_datos',
    python_callable=transformacion_datos,
    dag=dag
)

descarga_datos_task >> transformacion_datos_task

En este caso, "mi_libreria" es una biblioteca personalizada que contiene las funciones "descargar_datos" y "transformacion_datos". Estas funciones realizan las tareas necesarias para descargar los datos y transformarlos respectivamente.

La tarea "descarga_datos_task" llama a la función "descargar_datos" y la tarea "transformacion_datos_task" llama a la función "transformacion_datos". Al conectar estas tareas en nuestro flujo de trabajo, aseguramos que la tarea "transformacion_datos_task" no se ejecutará hasta que la tarea "descarga_datos_task" haya sido completada. De esta manera, garantizamos que los datos necesarios estarán disponibles para la siguiente tarea en el flujo de trabajo.

Aplicación práctica

Ejemplo práctico de creación de tarea personalizada y operador en Apache Airflow

Imaginemos que queremos crear una tarea que realice la descarga de archivos de un servidor FTP en una carpeta local. Para esto, crearemos una tarea personalizada llamada "DescargaFTPTask", que utilizará el operador "FTPDownloadOperator" que crearemos. Este operador será responsable de conectarse al servidor FTP y descargar los archivos.


from airflow.models import BaseOperator
from airflow.utils.decorators import apply_defaults
from ftplib import FTP

class DescargaFTPTask(BaseOperator):
    @apply_defaults
    def __init__(self, servidor_ftp, ruta_remota, ruta_local, *args, **kwargs):
        super(DescargaFTPTask, self).__init__(*args, **kwargs)
        self.servidor_ftp = servidor_ftp
        self.ruta_remota = ruta_remota
        self.ruta_local = ruta_local
    
    def execute(self, context):
        ftp = FTP(self.servidor_ftp)
        ftp.login()
        ftp.cwd(self.ruta_remota)
        archivos = ftp.nlst()
        for archivo in archivos:
            with open(f"{self.ruta_local}/{archivo}", "wb") as f:
                ftp.retrbinary(f"RETR {archivo}", f.write)

En este ejemplo, la clase "DescargaFTPTask" hereda de "BaseOperator" y define un método "execute" que realiza la conexión al servidor FTP, cambia a la ruta remota especificada, lista los archivos y los descarga a la carpeta local proporcionada.


from airflow.models import BaseOperator
from airflow.utils.decorators import apply_defaults
from ftplib import FTP

class FTPDownloadOperator(BaseOperator):
    @apply_defaults
    def __init__(self, server: str, remote_path: str, local_path: str, *args, **kwargs) -> None:
        super().__init__(*args, **kwargs)
        self.server = server
        self.remote_path = remote_path
        self.local_path = local_path
    
    def execute(self, context):
        ftp = FTP(self.server)
        ftp.login()
        ftp.cwd(self.remote_path)
        files = ftp.nlst()
        for file in files:
            with open(f"{self.local_path}/{file}", "wb") as f:
                ftp.retrbinary(f"RETR {file}", f.write)

En este segundo bloque de código, definimos el operador "FTPDownloadOperator", que también hereda de "BaseOperator". Este operador realiza la misma tarea que la tarea personalizada "DescargaFTPTask", pero encapsula la lógica de conexión FTP y descarga de archivos en un operador reutilizable.


from airflow.models import BaseOperator
from airflow.utils.decorators import apply_defaults
from ftplib import FTP
from airflow.operators.ftp_download_operator import FTPDownloadOperator

class DescargaFTPTask(BaseOperator):
    @apply_defaults
    def __init__(self, server, remote_path, local_path, *args, **kwargs):
        super().__init__(*args, **kwargs)
        self.server = server
        self.remote_path = remote_path
        self.local_path = local_path
    
    def execute(self, context):
        op = FTPDownloadOperator(
            task_id="descarga_ftp",
            server=self.server,
            remote_path=self.remote_path,
            local_path=self.local_path,
            dag=self.dag
        )
        op.execute(context)

Finalmente, para utilizar el operador "FTPDownloadOperator" en la tarea personalizada "DescargaFTPTask", creamos una instancia del operador dentro de la función "execute" de la tarea personalizada. De esta manera, la tarea personalizada delega la ejecución al operador para realizar la descarga de los archivos FTP.