ETL protsess andmete saamiseks e-postist Apache Airflow's.

ETL protsess andmete saamiseks e-postist Apache Airflow's.

Kuigi tehnoloogia areneb pidevalt, järgneb selle arengule alati hulk vananenud lähenemisviise. Seda võivad põhjustada sujuv üleminek, inimfaktor, tehnoloogilised vajadused või midagi muud. Andmete töötlemise valdkonnas on selles osas kõige silmatorkavamad andmeallikad. Kuigi me ei unista sellest lahti saada, saadetakse osa andmeid endiselt sõnumirakendustes ja e-kirjades, rääkimata veelgi arhailisematest formaatidest. Kutsun teid allpool vaatama ühte variante, kuidas Apache Airflow's e-kirjadest andmeid saada.

Eellugu

Paljusid andmeid edastatakse siiani e-posti teel, alates isiklike suhtluste lõpetamisest kuni ettevõtetevaheliste koostööstandarditeni. Hea, kui saad andmete hankimiseks kirjutada liidese või palgata inimesi, kes neid andmeid mugavamatesse allikatesse sisestaks, kuid sageli ei pruugi see olla teostatav. Konkreetselt, millega mina silmitsi seisin, oli tuntud CRM-süsteemi ühendamine andmehoidla ja hiljem OLAP-süsteemiga. Ajalooliselt on see süsteem olnud meie ettevõttele mugav teatud äri valdkonnas. Seetõttu soovisid kõik saada võimalust operaerida andmetega ka sellest kolmandast süsteemist. Esiteks uuriti loomulikult andmete saamise võimalust avatud API-st. Kahjuks ei katnud API kõiki vajalikke andmeid ja lihtsalt öeldes, oli see paljuski puudulik, samas ei soovinud tehniline tugi minna kaasa andmete täiendava funktsionaalsuse pakkumise osas. Sellegipoolest pakkus see süsteem võimalust perioodiliselt saada puuduvad andmed e-posti teel lingina arhiivi allalaadimiseks.

Tuleb märkida, et see ei olnud ainus juhtum, kus ettevõte soovis andmeid koguda e-kirjadest või sõnumirakendustest. Siiski ei saanud me selles olukorras mõjutada kolmandat ettevõtet, mis andmeid ainult selliselt edastas.

Apache Airflow

ETL-protsesside koostamisel kasutame me kõige sagedamini Apache Airflow'd. Et lugeja, kes ei ole selle tehnoloogiaga tuttav, paremini mõistaks, kuidas see kontekstis välja näeb, kirjeldan paar sissejuhatavat asja.

Apache Airflow on avatud platvorm, mida kasutatakse ETL (Extract-Transform-Loading) protsesside koostamiseks, täideviimiseks ja jälgimiseks Pythonis. Peamine mõisted Airflow's on suunatud suunamata graaf, kus graafi tipud on konkreetsed protsessid ja graafi servad on juhtimis- või info voog. Protsess võib lihtsalt kutsuda esile igasuguse Python'i funktsiooni või võib tal olla keerulisem loogika mitme funktsiooni järjestikuseks kutsumiseks klassi kontekstis. Kõige sagedasemate toimingute jaoks on juba palju valmis lahendusi, mida saab kasutada protsessidena. Nende hulka kuuluvad:

  • operaatorid — andmete edastamiseks ühte kohta teisest, näiteks DB tabelist andmehoidlasse;
  • sensorid — teatud sündmuse toimumise ootamiseks ja juhtimisvoo suunamiseks graafi järgmistele tippudele;
  • hookid — madalamal tasemel toimingute jaoks, näiteks andmete saamiseks DB tabelist (kasutatakse operaatorites);
  • jne.

Apache Airflow'i üksikasjalik kirjeldamine selles artiklis ei oleks otstarbekas. Lühikesi sissejuhatusi saab vaadata siin või siin.

Andmete saamise hook

Esiteks peab ülesande lahendamiseks kirjutama hooki, mille abil võiksime:

  • ühenduda e-posti kontoga;
  • leida vajaliku kirja;
  • saada andmeid kirjast.

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

class IMAPHook(BaseHook):
    def __init__(self, imap_conn_id):
        """
           IMAP hook, ettekandmiseks e-kirjade andmeid

           :param imap_conn_id:       E-posti ühenduse identifikaator
           :type imap_conn_id:        string
        """
        self.connection = self.get_connection(imap_conn_id)
        self.mail = None

    def authenticate(self):
        """ 
            Ühendame end e-postiga
        """
        mail = imaplib.IMAP4_SSL(self.connection.host)
        response, detail = mail.login(user=self.connection.login, password=self.connection.password)
        if response != "OK":
            raise AirflowException("Sisse logimine ebaõnnestus")
        else:
            self.mail = mail

    def get_last_mail(self, check_seen=True, box="INBOX", condition="(UNSEEN)"):
        """
            Meetod viimase e-kirja identifikaatori saamiseks,
            mis vastab otsingutingimustele

            :param check_seen:      Märgistada viimane kiri loetuks
            :type check_seen:       bool
            :param box:             Karbi nimi
            :type box:              string
            :param condition:       E-kirjade otsingutingimused
            :type condition:        string
        """
        self.authenticate()
        self.mail.select(mailbox=box)
        response, data = self.mail.search(None, condition)
        mail_ids = data[0].split()
        logging.info("Karbis leiti järgmised kirjad: " + str(mail_ids))

        if not mail_ids:
            logging.info("Uute kirjade leidmine ebaõnnestus")
            return None

        mail_id = mail_ids[0]

        # kui selliseid kirju on mitu
        if len(mail_ids) > 1:
            # märkame ülejäänud loetuks
            for id in mail_ids:
                self.mail.store(id, "+FLAGS", "\Seen")

            # tagastame viimase
            mail_id = mail_ids[-1]

        # kas tuleb viimane kiri märkida loetuks
        if not check_seen:
            self.mail.store(mail_id, "-FLAGS", "\Seen")

        return mail_id

Loogika on järgmine: ühendame end, leiame kõige värskema kirja, kui on ka teisi - ignoreerime neid. Kasutame just sellist funktsiooni, sest hilisemad kirjad sisaldavad kõiki varasemate andmeid. Kui see nii ei ole, siis võib tagastada kõikide kirjade massiivi või töödelda esimest, ning ülejäänud järgmise läbimise ajal. Üldiselt sõltub kõik nagu ikka ülesandest.

Lisame funktsiooni kahele abifunktsioonile: faili allalaadimiseks ja faili allalaadimiseks e-kirja lingilt. Pealegi saab need välja tuua operaatorisse, see sõltub kasutusastmest. Mida veel funktsiooni lisada, sõltub jällegi ülesandest: kui e-kirjas on kohe failid, siis saab allalaadida manuseid, kui andmed saadetakse e-kirjas, tuleb e-kiri parsida jne. Minu puhul saadetakse e-kirjas üks link arhiivile, mille pean panema kindlasse kohta ja käivitama edasise töötlemise protsessi.

    def download_from_url(self, url, path, chunk_size=128):
        """
            Faili allalaadimise meetod

            :param url:              Laadimisadress
            :type url:               string
            :param path:             Kuhu faili panna
            :type path:              string
            :param chunk_size:       Kui palju bite kirjutada
            :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):
        """
            Faili allalaadimise meetod e-kirja lingilt

            :param mail_id:         E-kirja identifikaator
            :type mail_id:          string
            :param path:            Kuhu faili panna
            :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)

Kood on lihtne, seega tõenäoliselt ei vajagi täiendavaid selgitusi. Räägin vaid maagilisest realt imap_conn_id. Apache Airflow salvestab ühendusparameetreid (kasutajanimi, parool, aadress ja muud parameetrid), millele saab juurdepääsu stringi identifikaatori kaudu. Visuaalselt näeb ühenduse haldamine välja nii

ETL protsess andmete saamiseks e-postist Apache Airflow's.

Andmete ootamise sensor

Kuna me juba oskame e-postiga ühendust luua ja andmeid saada, saame nüüd kirjutada sensori nende ootamiseks. Kirjutada kohe operaator, mis töötleb andmeid, kui need on olemas, ei õnnestunud, kuna saadud andmete põhjal töötavad ka teised protsessid, sealhulgas need, mis kasutavad seotud andmeid teistest allikatest (API, telefoniteenused, veebimetrika jne). Toon näite. CRM-süsteemis on ilmunud uus kasutaja, kuid me ei tea veel tema UUID-d. Seetõttu, kui proovime saada andmeid SIP-telefonist, saame kõnesid, mis on seotud tema UUID-ga, kuid ei suuda neid õigesti salvestada ja kasutada. Sellistes küsimustes on oluline arvestada andmete sõltuvust, eriti kui need on eri allikatest. Need on muidugi ebapiisavad meetmed andmete terviklikkuse säilitamiseks, kuid teatud juhtudel on need vajalikud. Ja ressursse tühjalt kulutada ei ole samuti mõistlik.

Seega käivitab meie sensor järgmised graafi tipud, kui e-kirjas on värsket teavet, ja märgib varasema teabe kehtetuks.

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

Saame ja kasutame andmeid

Andmete saamiseks ja töötlemiseks saab kirjutada eraldi operaatori või kasutada olemasolevaid. Kuna loogika on seni triviaalne — andmete hankimine e-kirjast, siis näitena pakun välja tavalise PythonOperatori

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)

# Tava tavaline graafi konfigureerimine
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",
)

# Määrame anduri
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",
)

# Funktsioon andmete saamiseks kirjast
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("Tühi postkast")

    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)

# Ülejäänud graafi sõlmede kirjeldus
...

# Seome graafi
mail_check_sensor >> prepare_mail_data
prepare_data >> ...
# Ülejäänud juhtimisvoogude kirjeldus

Muide, kui teie ettevõtte postkast on samuti mail.ru-s, siis ei ole teil võimalik otsida kirju teema, saatja jms järgi. Nad lubasid seda juba kauges 2016. aastal, kuid paistab, et on meelt muutnud. Lahendasin selle probleemi, luues vajalikele kirjadele eraldi kausta ja seadistades veebivaatlusesse filtriga vajalikud kirjad. Nii pääsevad sellesse kausta vaid vajalikud kirjad ja otsingutingimused on minu puhul lihtsalt (UNSEEN).

Kokkuvõttes on meil järgmine järjestus: kontrollime, kas on uusi kirju, mis vastavad tingimustele, ja kui on, siis laadime alla arhivi viite viimasest kirjast.
Viimaste ellipside taga on jäänud mainimata, et see arhiv lahti pakitakse, andmed arhivist puhastatakse ja töödeldakse ning lõpuks see kõik läheb edasi ETL protsessi konveierile, kuid see jääb artikli teema piiridest välja. Kui see oli huvitav ja kasulik, siis olen hea meelega valmis jätkama ETL lahenduste ja nende osade kirjeldamist Apache Airflow jaoks.

Allikas: habr.com

Osta usaldusväärne hostimine veebilehtede jaoks DDoS-i kaitsega, VPS VDS serverid 🔥 Osta usaldusväärne hostimine veebilehtede jaoks DDoS-i kaitsega, VPS VDS serverid | ProHoster