ETL procesi i marrjes së të dhënave nga email në Apache Airflow

ETL procesi i marrjes së të dhënave nga email në Apache Airflow

Pavarësisht nga përparimi i teknologjive, gjithmonë ka një sërë qasjeve të vjetra që i ndjekin ato. Kjo mund të jetë e lidhur me një kalim të ngadaltë, faktorë njerëzorë, nevoja teknologjike ose diçka tjetër. Në fushën e përpunimit të të dhënave, burimet e të dhënave janë ato që ilustrojnë më së miri këtë pjesë. Pavarësisht dëshirës sonë për t'u çliruar nga kjo, pjesa e të dhënave ende dërgohet përmes mesazheve dhe email-eve, pa përmendur format më archaike. Ju ftoj të shqyrtojmë nën këtë lidhje një nga mundësitë për Apache Airflow, që ilustron si mund të marrim të dhëna nga email-et.

Historia e mëparshme

ShumĂ« tĂ« dhĂ«na ende transmetohen pĂ«rmes email-it, duke filluar nga komunikimet ndĂ«rpersonale dhe pĂ«rfunduar me standardet e bashkĂ«punimit ndĂ«rmjet kompanive. ËshtĂ« mirĂ« nĂ«se arrijmĂ« tĂ« shkruajmĂ« njĂ« ndĂ«rfaqe pĂ«r marrjen e tĂ« dhĂ«nave apo tĂ« vendosim njerĂ«z nĂ« zyrĂ« qĂ« do tĂ« futnin kĂ«to informacione nĂ« burime mĂ« tĂ« pĂ«rshtatshme, por shpesh herĂ« njĂ« mundĂ«si e tillĂ« mund tĂ« mos ekzistojĂ«. Problemi konkret me tĂ« cilin u pĂ«rballa ishte lidhja e njohur tĂ« sistemit CRM me njĂ« depo tĂ« dhĂ«nash dhe mĂ« pas me sistemin OLAP. Historikisht, pĂ«rdorimi i kĂ«tij sistemi ishte i pĂ«rshtatshĂ«m pĂ«r njĂ« fushĂ« tĂ« caktuar biznesi. Prandaj, tĂ« gjithĂ« e donin mundĂ«sinĂ« pĂ«r tĂ« operuar me tĂ« dhĂ«nat dhe nga ky sistem i jashtĂ«m gjithashtu. Fillimisht, sigurisht, u shqyrtua mundĂ«sia e marrjes sĂ« tĂ« dhĂ«nave nga API i hapur. FatkeqĂ«sisht, API nuk mbulonte marrjen e tĂ« dhĂ«nave tĂ« nevojshme, dhe, pĂ«r ta thĂ«nĂ« thjesht, ishte nĂ« shumĂ« aspekte joefikas, dhe mbĂ«shtetje teknike nuk pranoi ose nuk mundi tĂ« ofronte njĂ« funksionalitet mĂ« tĂ« plotĂ«. MegjithatĂ«, ky sistem ofronte mundĂ«sinĂ« pĂ«r marrjen periodike tĂ« tĂ« dhĂ«nave tĂ« munguara nĂ« email nĂ« formĂ«n e njĂ« lidhjeje pĂ«r shkarkimin e arkivave.

Duhet theksuar se kjo nuk ishte rasti i vetëm për të cilin biznesi dëshironte të mbledhë të dhëna nga email-et ose mesazhet. Megjithatë, në këtë rast nuk mundëm të ndikojmë në kompaninë e tretë që ofronte disa të dhëna vetëm në këtë mënyrë.

Apache Airflow

Për ndërtimin e proceseve ETL, ne zakonisht përdorim Apache Airflow. Për të ndihmuar lexuesin që nuk është i njohur me këtë teknologji, do të përshkruaj disa hyrje, që ta kuptojë më mirë se si duket në kontekst dhe për gjithçka tjetër.

Apache Airflow është një platformë e lirë që përdoret për ndërtimin, ekzekutimin dhe monitorimin e proceseve ETL (Extract-Transform-Loading) në gjuhën Python. Koncepti kryesor në Airflow është grafin orientuar aciklik, ku kulmin e grafit përfaqësojnë proceset specifike dhe ëmbëlsirat e grafit janë rrjedha e kontrollit ose informacionit. Një proces mund të thërrasë thjesht çdo funksion Python, ose mund të ketë një logjikë më të komplikuar nga thirrjet e renditura të disa funksioneve në kontekstin e një klase. Për operacionet më të zakonshme, ekzistojnë tashmë shumë rreshta të gatshëm që mund të përdoren si procese. Këto përfshijnë:

  • operatoret — pĂ«r transferimin e tĂ« dhĂ«nave nga njĂ« vend nĂ« tjetrin, pĂ«r shembull nga njĂ« tabelĂ« DB nĂ« njĂ« depĂČ tĂ« dhĂ«nash;
  • sensort — pĂ«r tĂ« pritur qĂ« ndodhin ngjarje tĂ« caktuara dhe pĂ«r tĂ« drejtuar rrjedhĂ«n e kontrollit nĂ« kulminat e ardhshme tĂ« grafit;
  • hook — pĂ«r operacione mĂ« tĂ« ulta, pĂ«r shembull, pĂ«r tĂ« marrĂ« tĂ« dhĂ«na nga njĂ« tabelĂ« DB (pĂ«rdoren nĂ« operatoret);
  • etj.

Të përshkruash Apache Airflow në detaje në këtë artikull nuk do të kishte kuptim. Mund të shihni disa përmbledhje të shkurtra këtu ose këtu.

Hook për marrjen e të dhënave

Së pari, për të zgjidhur detyrën, duhet të shkruajmë një hook, me anë të cilit mund të:

  • lidhemi me email;
  • gjejmĂ« letrĂ«n e duhur;
  • marrim tĂ« dhĂ«na nga letra.

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

class IMAPHook(BaseHook):
    def __init__(self, imap_conn_id):
        """
           IMAP hook për marrjen e të dhënave nga emaili

           :param imap_conn_id:       Identifikuesi i lidhjes me email
           :type imap_conn_id:        string
        """
        self.connection = self.get_connection(imap_conn_id)
        self.mail = None

    def authenticate(self):
        """ 
            Këtu lidhemi me emailin
        """
        mail = imaplib.IMAP4_SSL(self.connection.host)
        response, detail = mail.login(user=self.connection.login, password=self.connection.password)
        if response != "OK":
            raise AirflowException("Hyrja dështoi")
        else:
            self.mail = mail

    def get_last_mail(self, check_seen=True, box="INBOX", condition="(UNSEEN)"):
        """
            Metoda për të marrë identifikuesin e emailit të fundit,
            që plotëson kushtet e kërkimit

            :param check_seen:      Shëno emailin e fundit si të lexuar
            :type check_seen:       bool
            :param box:             Emri i kutisë
            :type box:              string
            :param condition:       Kushtet e kërkimit për emaila
            :type condition:        string
        """
        self.authenticate()
        self.mail.select(mailbox=box)
        response, data = self.mail.search(None, condition)
        mail_ids = data[0].split()
        logging.info("Në kuti janë gjetur emailat si më poshtë: " + str(mail_ids))

        if not mail_ids:
            logging.info("Nuk u gjetën emaila të rinj")
            return None

        mail_id = mail_ids[0]

        # nëse ka disa nga këto emaila
        if len(mail_ids) > 1:
            # shëno të tjerët si të lexuar
            for id in mail_ids:
                self.mail.store(id, "+FLAGS", "\Seen")

            # kthe emailin e fundit
            mail_id = mail_ids[-1]

        # a duhet ta shënojmë emailin e fundit si të lexuar
        if not check_seen:
            self.mail.store(mail_id, "-FLAGS", "\Seen")

        return mail_id

Logjika Ă«shtĂ« kĂ«shtu: lidhemi, gjendim emailin e fundit mĂ« aktual, nĂ«se ka tĂ« tjerĂ« — i injorojmĂ« ata. PĂ«rdoret pikĂ«risht kjo funksion, sepse emailat mĂ« tĂ« vonshĂ«m pĂ«rmbajnĂ« tĂ« gjitha tĂ« dhĂ«nat e emailave tĂ« hershĂ«m. NĂ«se nuk Ă«shtĂ« kĂ«shtu, mund tĂ« kthejmĂ« njĂ« arrĂ« tĂ« tĂ« gjithĂ« emailave ose tĂ« pĂ«rpunojmĂ« tĂ« parin dhe tĂ« tjerĂ«t — nĂ« kalimin e ardhshĂ«m. NĂ« thelb, gjithçka varet si gjithmonĂ« nga detyra.

ShtojmĂ« dy funksione ndihmĂ«se nĂ« hook: pĂ«r shkarkimin e njĂ« skedari dhe pĂ«r shkarkimin e njĂ« skedari nga linku nĂ« email. MegjithatĂ«, ato mund tĂ« nxirren jashtĂ« nĂ« operator, kjo varet nga frekuenca e pĂ«rdorimit tĂ« kĂ«tij funksionaliteti. ÇfarĂ« tjetĂ«r tĂ« shtojmĂ« nĂ« hook, pĂ«rsĂ«ri varet nga detyra: nĂ«se nĂ« email vijnĂ« menjĂ«herĂ« skedarĂ«, atĂ«herĂ« mund tĂ« shkarkojmĂ« aplikacionet pĂ«r email, nĂ«se tĂ« dhĂ«nat vijnĂ« nĂ« email, atĂ«herĂ« duhet tĂ« pĂ«rpunojmĂ« emailin, etj. NĂ« rastin tim, emaili vjen me njĂ« lidhje pĂ«r njĂ« arkiv, tĂ« cilin duhet ta vendos nĂ« njĂ« vend tĂ« caktuar dhe tĂ« nis procesin e mĂ«tejshĂ«m tĂ« pĂ«rpunimit.

    def download_from_url(self, url, path, chunk_size=128):
        """
            Metoda për shkarkimin e një skedari

            :param url:              Adresa e shkarkimit
            :type url:               string
            :param path:             Ku të vendoset skedari
            :type path:              string
            :param chunk_size:       Sa shumë byte të shkruhen
            :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):
        """
            Metoda për shkarkimin e skedarit nga lidhja në email

            :param mail_id:         Identifikuesi i emailit
            :type mail_id:          string
            :param path:            Ku të vendoset skedari
            :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)

Kodi është i thjeshtë, për këtë arsye vështirë se ka nevojë për shpjegime shtesë. Po flas vetëm për rreshtin magjik imap_conn_id. Apache Airflow ruan parametrat e lidhjeve (emri i përdoruesit, fjalëkalimi, adresa dhe parametrat e tjerë), të cilave mund t'u qaseni përmes një identifikuesi të rreshtit. Visualisht, menaxhimi i lidhjeve duket kështu

ETL procesi i marrjes së të dhënave nga email në Apache Airflow

Sensor për pritjen e të dhënave

Tani se ne nevojitet të lidhemi dhe të marrim të dhëna nga posta, tani mund të shkruajmë një sensor për t'i pritur ato. Nuk arrita të shkruaj menjëherë një operator që do të përpunonte të dhënat, në rastin tim, sepse proceset e tjera punojnë gjithashtu mbi të dhënat e marra nga posta, përfshirë ato që marrin të dhëna të lidhura nga burime të tjera (API, telefonia, metrikat e webit, etj.). Do të jap një shembull. Në sistemin CRM ka një përdorues të ri, dhe ne ende nuk e dimë për UUID-në e tij. Atëherë, kur përpiqemi të marrim të dhëna nga telefonia SIP, ne do të marrim thirrjet e lidhura me UUID-në e tij, por nuk do të jemi në gjendje t'i ruajmë dhe t'i përdorim ato siç duhet. Në këto çështje është e rëndësishme të kemi parasysh varësinë e të dhënave, veçanërisht nëse ato vijnë nga burime të ndryshme. Këto, sigurisht, janë masa të pamjaftueshme për të ruajtur integritetin e të dhënave, por në disa raste janë të nevojshme. Po ashtu, të zënë burimet kot është edhe irracional.

Në këtë mënyrë, sensori ynë do të aktivizojë majat e mëpasshme të grafit, nëse ka informacione të freskëta në postë, si dhe do të shënojë informacionin e mëparshëm si të pavlefshëm.

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

Marrim dhe përdorim të dhëna

PĂ«r marrjen dhe pĂ«rpunimin e tĂ« dhĂ«nave mund tĂ« shkruhet njĂ« operator i veçantĂ«, ose mund tĂ« pĂ«rdoren ato tĂ« gatshme. Duke qenĂ« se logjika pĂ«r momentin Ă«shtĂ« triviale — tĂ« marrim tĂ« dhĂ«na nga emaili, pĂ«r shembull propozoj njĂ« PythonOperator standard.

nga airflow.models import DAG

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

start_date = datetime(2020, 4, 4)

# Konfigurimi standard i DAG-ut
args = {
    "owner": "shembull",
    "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",
)

# Përcaktojmë sensorin
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",
)

# Funksioni për të marrë të dhëna nga email-i
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("Mailbox i zbrazët")

    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)

# Përshkrimi i kulmave të tjera të DAG-ut
...

# Caktojmë lidhjen në DAG
mail_check_sensor >> prepare_mail_data
prepare_data >> ...
# Përshkrimi i kontrollit të tjera

Për më tepër, nëse posti juaj korporativ është gjithashtu në mail.ru, atëherë nuk do t'ju jetë e mundur të kërkoni mesazhe sipas temës, dërguesit etj. Ata premtuan të implementojnë këtë në vitin 2016, por duket se e kanë revokuar këtë vendim. Unë e zgjodha këtë problem duke krijuar një dosje të veçantë për email-et që më duhen dhe duke konfiguruar një filtër në ndërfaqen e uebit për email-et e nevojshme. Kështu, në këtë dosje i bien vetëm email-et e kërkuara dhe kushti për kërkimin në rastin tim është thjesht (UNSEEN).

Në përmbledhje, ne kemi këtë renditje: kontrollojmë nëse ka email-e të reja që përmbushin kushtet, nëse ka, atëherë shkarkojmë arkivin nga linku në email-in e fundit.
Nën shumëpikëshat e fundit është lënë pa përmendur se ky arkiv do të shpaketohet, të dhënat nga arkivi do të pastrohen dhe përpunohen, dhe në fund gjithçka do të shkojë përpara në procesin ETL, por kjo është jashtë temës së këtij artikulli. Nëse u duk e interesuar dhe e dobishme, me kënaqësi do të vazhdoj të përshkruaj zgjidhjet ETL dhe pjesët e tyre për Apache Airflow.

Burimi: habr.com

Blini hosting tĂ« besueshĂ«m pĂ«r faqe interneti me mbrojtje nga DDoS, serverĂ« VPS VDS đŸ”„ Blini hosting tĂ« besueshĂ«m pĂ«r faqe interneti me mbrojtje nga DDoS, serverĂ« VPS VDS | ProHoster