Proceso ETL de obtención de datos desde el correo electrónico en Apache Airflow

Proceso ETL de obtención de datos desde el correo electrónico en Apache Airflow

Por mucho que las tecnologías avancen, siempre hay una serie de enfoques obsoletos que los acompañan. Esto puede deberse a una transición gradual, al factor humano, a necesidades tecnológicas o a otras razones. En el ámbito del procesamiento de datos, los más representativos son las fuentes de datos. Por más que soñemos con deshacernos de ello, mientras parte de los datos se envíe a través de mensajeros y correos electrónicos, sin mencionar formatos aún más arcaicos. Te invito a desglosar una de las variantes para Apache Airflow, que ilustra cómo se pueden extraer datos de correos electrónicos.

Antecedentes

Muchos datos aún se transmiten por correo electrónico, desde las comunicaciones interpersonales hasta los estándares de interacción entre empresas. Es bueno si se puede crear una interfaz para obtener datos o si se pueden sentar en la oficina a las personas que ingresen esta información en fuentes más convenientes, pero a menudo esta posibilidad simplemente no está disponible. El desafío específico con el que me encontré fue conectar un conocido sistema de CRM a un almacenamiento de datos y, posteriormente, a un sistema OLAP. Históricamente, en nuestra empresa, utilizar este sistema ha sido conveniente en un área de negocio en particular. Por ello, todos querían poder operar con datos de este sistema externo también. En primer lugar, se examinó la posibilidad de obtener datos a través de una API abierta. Desafortunadamente, la API no cubría la obtención de todos los datos necesarios y, siendo claros, tenía muchas limitaciones, además de que el soporte técnico no quiso o no pudo colaborar para ofrecer una funcionalidad más completa. Sin embargo, este sistema ofrecía la posibilidad de recibir periódicamente los datos faltantes por correo electrónico en forma de un enlace para descargar el archivo.

Es importante señalar que este no fue el único caso en el que el negocio quería recopilar datos de correos electrónicos o mensajeros. Sin embargo, en este caso no pudimos influir en la empresa externa que proporciona parte de los datos solo de esta manera.

Apache Airflow

Para construir procesos ETL, usamos con mayor frecuencia Apache Airflow. Para ayudar al lector, que no está familiarizado con esta tecnología, a entender mejor cómo se presenta en este contexto y en general, describiré un par de introducciones.

Apache Airflow es una plataforma abierta que se utiliza para construir, ejecutar y monitorear procesos ETL (Extract-Transform-Loading) en Python. El concepto principal en Airflow es el grafo dirigido acíclico, donde los nodos del grafo son procesos específicos, y las aristas del grafo representan el flujo de control o información. Un proceso puede simplemente invocar cualquier función de Python, o puede tener una lógica más compleja que consiste en la llamada secuencial de varias funciones en el contexto de una clase. Para las operaciones más comunes, ya existen numerosas implementaciones listas que se pueden usar como procesos. Entre estas implementaciones se encuentran:

  • operadores: para mover datos de un lugar a otro, por ejemplo, de una tabla de base de datos a un almacén de datos;
  • sensores: para esperar la ocurrencia de determinado evento y dirigir el flujo de control a los siguientes nodos del grafo;
  • ganchos: para operaciones de menor nivel, por ejemplo, para obtener datos de una tabla de base de datos (se utilizan en los operadores);
  • etc.

No sería sensato describir Apache Airflow en detalle en este artículo. Se pueden ver introducciones breves en aquí o aquí.

Gancho para obtener datos

Primero que nada, para resolver la tarea, necesitamos escribir un gancho que nos permita:

  • conectarnos al correo electrónico;
  • encontrar el correo necesario;
  • obtener datos del correo.

from airflow.hooks.base_hook import BaseHook
import imaplib
import logging

class IMAPHook(BaseHook):
    def __init__(self, imap_conn_id):
        """
           IMAP hook para obtener datos del correo electrónico

           :param imap_conn_id:       Identificador de conexión con el correo
           :type imap_conn_id:        string
        """
        self.connection = self.get_connection(imap_conn_id)
        self.mail = None

    def authenticate(self):
        """ 
            Conectamos al correo
        """
        mail = imaplib.IMAP4_SSL(self.connection.host)
        response, detail = mail.login(user=self.connection.login, password=self.connection.password)
        if response != "OK":
            raise AirflowException("Error de inicio de sesión")
        else:
            self.mail = mail

    def get_last_mail(self, check_seen=True, box="INBOX", condition="(UNSEEN)"):
        """
            Método para obtener el identificador del último correo, 
            que cumple con las condiciones de búsqueda

            :param check_seen:      Marcar el último correo como leído
            :type check_seen:       bool
            :param box:             Nombre de la bandeja
            :type box:              string
            :param condition:       Condiciones para buscar correos
            :type condition:        string
        """
        self.authenticate()
        self.mail.select(mailbox=box)
        response, data = self.mail.search(None, condition)
        mail_ids = data[0].split()
        logging.info("Se encontraron los siguientes correos en la bandeja: " + str(mail_ids))

        if not mail_ids:
            logging.info("No se encontraron nuevos correos")
            return None

        mail_id = mail_ids[0]

        # si hay varios correos
        if len(mail_ids) > 1:
            # marcamos los demás como leídos
            for id in mail_ids:
                self.mail.store(id, "+FLAGS", "\Seen")

            # retornamos el último
            mail_id = mail_ids[-1]

        # se debe marcar el último como leído?
        if not check_seen:
            self.mail.store(mail_id, "-FLAGS", "\Seen")

        return mail_id

La lógica es la siguiente: nos conectamos, encontramos el correo más reciente, si hay otros, los ignoramos. Se utiliza esta función porque los correos más recientes contienen todos los datos de los anteriores. Si no es así, se puede devolver un arreglo con todos los correos o procesar el primero, y los demás en la siguiente pasada. En general, todo depende de la tarea.

Añadimos al gancho dos funciones auxiliares: una para descargar archivos y otra para descargar un archivo a través de un enlace en el correo. Por cierto, se pueden separar en un operador, dependiendo de la frecuencia de uso de esta funcionalidad. ¿Qué más agregar al gancho? Nuevamente, depende de la tarea: si el correo llega con archivos adjuntos, se pueden descargar estos adjuntos; si los datos llegan en el correo, entonces es necesario analizar el correo, etc. En mi caso, el correo llega con un enlace a un archivo comprimido que necesito colocar en un lugar específico y lanzar el proceso de procesamiento posterior.

    def download_from_url(self, url, path, chunk_size=128):
        """
            Método para descargar un archivo

            :param url:              Dirección de descarga
            :type url:               string
            :param path:             Dónde colocar el archivo
            :type path:              string
            :param chunk_size:       Cuántos bytes escribir
            :type chunk_size:        int
        """
        r = requests.get(url, stream=True)
        with open(path, "wb") as fd:
            for chunk in r.iter_content(chunk_size=chunk_size):
                fd.write(chunk)

    def download_mail_href_attachment(self, mail_id, path):
        """
            Método para descargar un archivo a través de un enlace en el correo

            :param mail_id:         Identificador del correo
            :type mail_id:          string
            :param path:            Dónde colocar el archivo
            :type path:             string
        """
        response, data = self.mail.fetch(mail_id, "(RFC822)")
        raw_email = data[0][1]
        raw_soup = raw_email.decode().replace("r", "").replace("n", "")
        parse_soup = BeautifulSoup(raw_soup, "html.parser")
        link_text = ""

        for a in parse_soup.find_all("a", href=True, text=True):
            link_text = a["href"]

        self.download_from_url(link_text, path)

El código es simple, por lo que probablemente no necesite más explicaciones. Solo comentaré sobre la línea mágica imap_conn_id. Apache Airflow almacena los parámetros de conexión (nombre de usuario, contraseña, dirección y otros parámetros), los cuales se pueden acceder mediante un identificador en forma de cadena. Visualmente, la gestión de conexiones se ve así.

Proceso ETL de obtención de datos desde el correo electrónico en Apache Airflow

Sensor para esperar datos

Dado que ya sabemos cómo conectarnos y obtener datos del correo, ahora podemos escribir un sensor para esperarlos. No pude escribir al mismo tiempo un operador que procese los datos, si están disponibles, ya que los datos obtenidos del correo son parte de otros procesos, incluidos aquellos que obtienen datos relacionados de otras fuentes (API, telefonía, métricas web, etc.). Pondré un ejemplo. En el sistema CRM apareció un nuevo usuario, y aún no sabemos su UUID. Entonces, al intentar obtener datos de la telefonía SIP, recibiremos llamadas vinculadas a su UUID, pero no podremos almacenarlas y utilizarlas correctamente. En estos asuntos, es importante tener en cuenta la dependencia de los datos, especialmente si provienen de diferentes fuentes. Esto, por supuesto, son medidas insuficientes para mantener la integridad de los datos, pero en algunos casos son necesarias. Y también es irracional ocupar recursos en vano.

Así, nuestro sensor activará las siguientes cumbres del grafo si hay información nueva en el correo, y también marcará como obsoleta la información anterior.

from airflow.sensors.base_sensor_operator import BaseSensorOperator
from airflow.utils.decorators import apply_defaults
from my_plugin.hooks.imap_hook import IMAPHook

class MailSensor(BaseSensorOperator):
    @apply_defaults
    def __init__(self, conn_id, check_seen=True, box="Inbox", condition="(UNSEEN)", *args, **kwargs):
        super().__init__(*args, **kwargs)
        self.conn_id = conn_id
        self.check_seen = check_seen
        self.box = box
        self.condition = condition

    def poke(self, context):
        conn = IMAPHook(self.conn_id)
        mail_id = conn.get_last_mail(check_seen=self.check_seen, box=self.box, condition=self.condition)

        if mail_id is None:
            return False
        else:
            return True

Obtenemos y utilizamos datos

Para obtener y procesar datos, se puede escribir un operador separado o utilizar uno ya existente. Dado que por ahora la lógica es trivial — obtener datos del correo, propongo usar el operador estándar PythonOperator como ejemplo.

from airflow.models import DAG

from airflow.operators.python_operator import PythonOperator
from airflow.sensors.my_plugin import MailSensor
from my_plugin.hooks.imap_hook import IMAPHook

start_date = datetime(2020, 4, 4)

# Configuración estándar del gráfico
args = {
    "owner": "example",
    "start_date": start_date,
    "email": ["home@home.ru"],
    "email_on_failure": False,
    "email_on_retry": False,
    "retry_delay": timedelta(minutes=15),
    "provide_context": False,
}

dag = DAG(
    dag_id="test_etl",
    default_args=args,
    schedule_interval="@hourly",
)

# Definimos el sensor
mail_check_sensor = MailSensor(
    task_id="check_new_emails",
    poke_interval=10,
    conn_id="mail_conn_id",
    timeout=10,
    soft_fail=True,
    box="my_box",
    dag=dag,
    mode="poke",
)

# Función para obtener datos del correo
def prepare_mail():
    imap_hook = IMAPHook("mail_conn_id")
    mail_id = imap_hook.get_last_mail(check_seen=True, box="my_box")
    if mail_id is None:
        raise AirflowException("Bandeja de entrada vacía")

    conn.download_mail_href_attachment(mail_id, "./path.zip")

prepare_mail_data = PythonOperator(task_id="prepare_mail_data", default_args=args, dag=dag, python_callable= prepare_mail)

# Descripción de otros nodos del gráfico
...

# Establecemos la relación en el gráfico
mail_check_sensor >> prepare_mail_data
prepare_data >> ...
# Descripción de otros flujos de control

Por cierto, si tu correo corporativo también está en mail.ru, no tendrás acceso a la búsqueda de correos por tema, remitente, etc. Ellos prometieron hacerlo en 2016, pero, aparentemente, se echaron atrás. Yo solucioné este problema creando una carpeta separada para los correos que necesitaba y configurando un filtro en la interfaz web del correo para esos correos. Así, solo los correos necesarios van a esa carpeta y las condiciones para la búsqueda en mi caso son simplemente (UNSEEN).

En resumen, tenemos la siguiente secuencia: verificamos si hay nuevos correos que cumplan las condiciones, si los hay, descargamos el archivo comprimido del enlace en el último correo.
Bajo los últimos puntos suspensivos se omite que este archivo comprimido será descomprimido, los datos del archivo serán limpiados y procesados, y al final todo esto se enviará al flujo del proceso ETL, pero eso ya se sale del tema del artículo. Si te parece interesante y útil, con gusto continuaré describiendo soluciones ETL y sus partes para Apache Airflow.

Fuente: habr.com

Compra un hosting fiable para sitios web con protección contra DDoS, servidores VPS VDS 🔥 Compra un hosting fiable para sitios web con protección contra DDoS, servidores VPS VDS | ProHoster