ETL-proces voor het verkrijgen van gegevens uit e-mail in Apache Airflow

ETL-proces voor het verkrijgen van gegevens uit e-mail in Apache Airflow

Hoezeer technologieën zich ook ontwikkelen, er is altijd een reeks verouderde benaderingen die met de ontwikkeling meegaan. Dit kan komen door een geleidelijke overgang, de menselijke factor, technologische behoeften of iets anders. In de wereld van gegevensverwerking zijn de dataleveranciers het meest kenmerkend in dit opzicht. Hoezeer we ook zouden willen ontsnappen aan deze situatie, zolang een deel van de gegevens via berichtenapps en e-mails wordt verzonden, om nog maar te zwijgen van de meer archaïsche formaten. Ik nodig je uit om onder deze regel een van de opties voor Apache Airflow te bespreken, die illustreert hoe je gegevens uit e-mails kunt ophalen.

Achtergrond

Veel gegevens worden nog steeds via e-mail verzonden, variƫrend van interpersoonlijke communicatie tot standaarden voor samenwerking tussen bedrijven. Het is goed als je een interface kunt schrijven om gegevens te ontvangen of mensen in het kantoor kunt hebben die deze informatie in meer toegankelijke bronnen invoeren, maar vaak is die mogelijkheid er gewoon niet. De specifieke uitdaging waarmee ik werd geconfronteerd, was het verbinden van een bekende CRM-systeem met een datawarehouse en vervolgens met een OLAP-systeem. Historisch gezien was het gebruik van dit systeem in een bepaald zakelijke domein handig voor ons bedrijf. Daarom wilde iedereen de mogelijkheid hebben om met gegevens uit dit externe systeem te werken. Uiteraard werd in eerste instantie gekeken naar de mogelijkheden om gegevens uit de open API te verkrijgen. Helaas dekte de API het verkrijgen van alle noodzakelijke gegevens niet en, eenvoudig gezegd, was het in veel opzichten krom, en de technische ondersteuning weigerde of kon niet tegemoetkomen door meer uitgebreide functionaliteit te bieden. Deze systeem bood echter de mogelijkheid om periodiek ontbrekende gegevens per e-mail te ontvangen in de vorm van een link voor het downloaden van een archief.

Het moet worden opgemerkt dat dit niet de enige casus was waarbij het bedrijf gegevens uit e-mails of berichtenapps wilde verzamelen. In dit geval konden we echter geen invloed uitoefenen op het externe bedrijf dat een deel van de gegevens alleen op deze manier verstrekte.

Apache Airflow

Voor het opzetten van ETL-processen gebruiken we vaak Apache Airflow. Om de lezer, die niet bekend is met deze technologie, beter te laten begrijpen hoe dit eruitziet in context en in het algemeen, beschrijf ik een paar inleidingen.

Apache Airflow is een gratis platform dat wordt gebruikt voor het opbouwen, uitvoeren en monitoren van ETL (Extract-Transform-Load) processen in Python. Het belangrijkste concept in Airflow is de gericht acyclische grafiek, waarbij de knopen van de grafiek concrete processen zijn en de verbindingen de stroom van controle of informatie vertegenwoordigen. Een proces kan eenvoudigweg een Python-functie aanroepen, maar kan ook een complexere logica hebben van opeenvolgende aanroepen van verschillende functies in de context van een klasse. Voor de meest voorkomende operaties zijn er al vele kant-en-klare modules die als processen kunnen worden gebruikt. Onder deze modules vallen:

  • operatoren – voor het verplaatsen van gegevens van de ene plaats naar de andere, bijvoorbeeld van een database tabel naar een datawarehouse;
  • sensoren – voor het wachten op het optreden van een bepaald evenement en het doorsturen van de controleflow naar de volgende knopen in de grafiek;
  • hooks – voor lager niveau operaties, bijvoorbeeld voor het ophalen van gegevens uit een database tabel (worden gebruikt in operatoren);
  • enzovoorts.

Het zou niet doelmatig zijn om Apache Airflow uitgebreid in dit artikel te beschrijven. Korte inleidingen zijn te vinden op hier of hier.

Hook voor gegevensophaling

In de eerste plaats moet er een hook worden geschreven waarmee we kunnen:

  • verbinden met e-mail;
  • de juiste e-mail vinden;
  • gegevens uit de e-mail ophalen.

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

class IMAPHook(BaseHook):
    def __init__(self, imap_conn_id):
        """
           IMAP-hook voor het ophalen van gegevens van e-mail

           :param imap_conn_id:       E-mail connectie ID
           :type imap_conn_id:        string
        """
        self.connection = self.get_connection(imap_conn_id)
        self.mail = None

    def authenticate(self):
        """ 
            Verbinden met de e-mail
        """
        mail = imaplib.IMAP4_SSL(self.connection.host)
        response, detail = mail.login(user=self.connection.login, password=self.connection.password)
        if response != "OK":
            raise AirflowException("Inloggen mislukt")
        else:
            self.mail = mail

    def get_last_mail(self, check_seen=True, box="INBOX", condition="(UNSEEN)"):
        """
            Methode om de ID van de laatste e-mail te verkrijgen, 
            voldoe aan de zoekcriteria

            :param check_seen:      Markeer de laatste e-mail als gelezen
            :type check_seen:       bool
            :param box:             Naam van de mailbox
            :type box:              string
            :param condition:       Zoekcriteria voor e-mails
            :type condition:        string
        """
        self.authenticate()
        self.mail.select(mailbox=box)
        response, data = self.mail.search(None, condition)
        mail_ids = data[0].split()
        logging.info("In de mailbox zijn de volgende e-mails gevonden: " + str(mail_ids))

        if not mail_ids:
            logging.info("Geen nieuwe e-mails gevonden")
            return None

        mail_id = mail_ids[0]

        # als er meerdere e-mails zijn
        if len(mail_ids) > 1:
            # markeer de overige als gelezen
            for id in mail_ids:
                self.mail.store(id, "+FLAGS", "\Seen")

            # retourneer de laatste
            mail_id = mail_ids[-1]

        # moet de laatste als gelezen worden gemarkeerd
        if not check_seen:
            self.mail.store(mail_id, "-FLAGS", "\Seen")

        return mail_id

De logica is als volgt: we verbinden, vinden de meest actuele e-mail; als er andere zijn, negeren we die. Deze functie wordt precies zo gebruikt omdat latere e-mails alle gegevens van eerdere bevatten. Als dat niet het geval is, kan een array van alle e-mails worden teruggegeven of de eerste worden verwerkt en de overige bij de volgende doorloop. Kortom, zoals altijd hangt het allemaal af van de taak.

We are adding two helper functions to the hook: one for downloading a file and one for downloading a file via a link from an email. By the way, they can be extracted to the operator, depending on how often this functionality is used. What else to add to the hook again depends on the task: if the email contains files right away, you can download the attachments to the email; if data comes in the email, you need to parse the email, etc. In my case, the email comes with a single link to an archive that I need to put in a specific location and start the further processing.

    def download_from_url(self, url, path, chunk_size=128):
        """
            Method for downloading a file

            :param url:              Download address
            :type url:               string
            :param path:             Where to place the file
            :type path:              string
            :param chunk_size:       Number of bytes to write at once
            :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):
        """
            Method for downloading a file via a link from an email

            :param mail_id:         Identifier of the email
            :type mail_id:          string
            :param path:            Where to place the file
            :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)

The code is simple, so it hardly needs further explanation. I will only mention the magic line imap_conn_id. Apache Airflow stores connection parameters (login, password, address, and other settings) that can be accessed by a string identifier. Visually, managing connections looks like this.

ETL-proces voor het verkrijgen van gegevens uit e-mail in Apache Airflow

Sensor for waiting for data

Aangezien we al weten hoe we verbinding kunnen maken en gegevens uit e-mail kunnen ophalen, kunnen we nu een sensor schrijven om ze te wachten. Het lukte me niet om onmiddellijk een operator te schrijven die de gegevens zou verwerken, als die er zijn, omdat op basis van de ontvangen gegevens uit de e-mail ook andere processen werken, waaronder die die gekoppelde gegevens uit andere bronnen halen (API, telefonie, webstatistieken, etc.). Ik geef een voorbeeld. In het CRM-systeem is er een nieuwe gebruiker verschenen, en we weten nog niet wat zijn UUID is. Bij het proberen om gegevens van SIP-telefonie te krijgen, zullen we gesprekken krijgen die aan zijn UUID zijn gekoppeld, maar we zullen ze niet correct kunnen opslaan of gebruiken. Bij dergelijke vragen is het belangrijk om de afhankelijkheid van gegevens in gedachten te houden, vooral als ze uit verschillende bronnen komen. Dit zijn natuurlijk niet voldoende maatregelen om de integriteit van gegevens te waarborgen, maar soms zijn ze noodzakelijk. En ook het ongebruikt bezet houden van middelen is niet rationeel.

Zo zal onze sensor de volgende knooppunten van de grafiek activeren als er nieuwe informatie in de e-mail is, en ook de vorige informatie als verouderd markeren.

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

Gegevens ophalen en gebruiken

Voor het verkrijgen en verwerken van gegevens kan een aparte operator worden geschreven, of een bestaande worden gebruikt. Aangezien de logica voorlopig triviaal is — gegevens uit een e-mail ophalen, stel ik voor om de standaard PythonOperator als voorbeeld te nemen.

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)

# Standaard configuratie van de grafiek
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",
)

# Sensor definitie
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",
)

# Functie om gegevens uit de e-mail te halen
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("Lege mailbox")

    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)

# Beschrijving van andere knooppunten in de grafiek
...

# Verbinding in de grafiek instellen
mail_check_sensor >> prepare_mail_data
prepare_data >> ...
# Beschrijving van andere controlestromen

Trouwens, als je zakelijke mail ook op mail.ru staat, dan kun je geen e-mails zoeken op onderwerp, afzender, enz. Ze beloofden dit in 2016, maar hebben blijkbaar van gedachten veranderd. Ik heb dit probleem opgelost door voor de benodigde e-mails een aparte map te maken en in de webinterface van de mail een filter in te stellen voor de gewenste e-mails. Zo komen alleen de benodigde e-mails in deze map en zijn de zoekvoorwaarden in mijn geval eenvoudig (UNSEEN).

Samenvattend hebben we de volgende volgorde: we controleren of er nieuwe e-mails zijn die aan de voorwaarden voldoen, en als dat zo is, downloaden we het archief via de link in de laatste e-mail.
Onder de laatste uitroeptekens wordt weggelaten dat dit archief zal worden uitgepakt, de gegevens uit het archief zullen worden schoongemaakt en verwerkt, en uiteindelijk zal dit alles verder gaan naar de ETL-proceslijn, maar dit valt al buiten het onderwerp van het artikel. Als je het interessant en nuttig vond, dan beschrijf ik graag verdere ETL-oplossingen en hun onderdelen voor Apache Airflow.

Bron: habr.com

Koop betrouwbare webhosting met bescherming tegen DDoS, VPS VDS servers šŸ”„ Koop betrouwbare webhosting met bescherming tegen DDoS, VPS VDS servers | ProHoster