
Indiferent cât de mult progresează tehnologia, în urma acesteia există întotdeauna o serie de metode învechite. Acest lucru poate fi determinat de tranziția lină, factorul uman, necesitățile tehnologice sau alte motive. În domeniul procesării datelor, cele mai relevante exemple în acest sens sunt sursele de date. Indiferent de cât de mult ne dorim să scăpăm de acestea, o parte din date continuă să fie transmise prin mesaje și emailuri, fără a menționa formatele mai arhaice. Vă invit să discutăm subiectul unui model pentru Apache Airflow, ilustrând cum se pot obține date din emailuri.
Povestea
Multe date sunt încă transmise prin email, începând cu comunicările interumane și terminând cu standardele de interacțiune între companii. Este ideal dacă se reușește să se scrie o interfață pentru obținerea datelor sau să se angajeze oameni în birou care să introducă aceste informații în surse mai accesibile, dar de multe ori această posibilitate nu există. Problema specifică cu care m-am confruntat a fost integrarea unei cunoscute sisteme CRM cu un depozit de date, iar ulterior — cu un sistem OLAP. Istoric, utilizarea acestui sistem a fost convenabilă pentru compania noastră într-un domeniu specific de afaceri. Prin urmare, tuturor le-ar fi plăcut să aibă posibilitatea de a opera cu datele din acest sistem terț. În prima fază, desigur, am studiat posibilitatea obținerii de date din API-ul deschis. Din păcate, API-ul nu acoperea obținerea tuturor datelor necesare și, simplu spus, era în mare parte imperfect, iar suportul tehnic nu a dorit sau nu a putut să se adapteze pentru a oferi un funcțional suplimentar mai cuprinzător. Totuși, acest sistem oferea posibilitatea de a primi periodic datele lipsă pe email sub formă de link pentru descărcarea arhivelor.
Este de remarcat că acesta nu a fost singurul caz în care afacerea a dorit să colecteze date din emailuri sau mesagerii. Cu toate acestea, în acest caz nu am putut influența compania terță, care furniza o parte din date doar în acest mod.
Apache Airflow
Pentru a construi procese ETL, folosim cel mai adesea Apache Airflow. Pentru ca cititorul, care nu este familiarizat cu această tehnologie, să înțeleagă mai bine cum arată aceasta în context și în general, voi descrie câteva aspecte introductive.
Apache Airflow este o platformă liberă care este utilizată pentru a construi, executa și monitoriza procese ETL (Extract-Transform-Loading) în limbajul Python. Conceptul principal în Airflow este un grafic orientat aciclic, unde nodurile graficului sunt procesele specifice, iar muchiile graficului reprezintă fluxul de control sau informație. Un proces poate apela pur și simplu orice funcție Python sau poate avea o logică mai complexă de apeluri secvențiale ale mai multor funcții în contextul unei clase. Pentru cele mai frecvente operații, există deja numeroase soluții gata făcute care pot fi utilizate ca procese. Printre aceste soluții se numără:
- operatoare - pentru transferul datelor dintr-un loc în altul, de exemplu dintr-un tabel BD în depozitul de date;
- senzori - pentru a aștepta apariția unui anumit eveniment și pentru a direcționa fluxul de control către următoarele noduri ale graficului;
- hook-uri - pentru operațiuni de nivel inferior, de exemplu, pentru a obține date dintr-un tabel BD (folosite în operatori);
- etc.
A descrie Apache Airflow în detaliu în acest articol nu ar fi eficient. Puteți viziona introducerile scurte în sau .
Hook pentru obținerea datelor
În primul rând, pentru a rezolva problema, trebuie să scriem un hook, cu ajutorul căruia am putea:
- conecta la e-mail;
- găsi e-mailul dorit;
- obține datele din e-mail.
from airflow.hooks.base_hook import BaseHook
import imaplib
import logging
class IMAPHook(BaseHook):
def __init__(self, imap_conn_id):
"""
IMAP hook pentru a obține date de la e-mail
:param imap_conn_id: Identificatorul conexiunii la e-mail
:type imap_conn_id: string
"""
self.connection = self.get_connection(imap_conn_id)
self.mail = None
def authenticate(self):
"""
Ne conectăm la 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("Conectarea a eșuat")
else:
self.mail = mail
def get_last_mail(self, check_seen=True, box="INBOX", condition="(UNSEEN)"):
"""
Metodă pentru a obține identificatorul ultimului e-mail,
care respectă condițiile de căutare
:param check_seen: Marchează ultimul e-mail ca citit
:type check_seen: bool
:param box: Numele cutiei poștale
:type box: string
:param condition: Condițiile de căutare a e-mailurilor
: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 cutie au fost găsite următoarele e-mailuri: " + str(mail_ids))
if not mail_ids:
logging.info("Nu au fost găsite e-mailuri noi")
return None
mail_id = mail_ids[0]
# dacă există mai multe astfel de e-mailuri
if len(mail_ids) > 1:
# marcăm celelalte ca fiind citite
for id in mail_ids:
self.mail.store(id, "+FLAGS", "\Seen")
# returnăm ultimul
mail_id = mail_ids[-1]
# trebuie să marcăm ultimul ca citit
if not check_seen:
self.mail.store(mail_id, "-FLAGS", "\Seen")
return mail_idLogica este următoarea: ne conectăm, găsim ultimul e-mail cel mai relevant, iar dacă există altele — le ignorăm. Se folosește exact această funcție, deoarece e-mailurile mai recente conțin toate datele e-mailurilor anterioare. Dacă nu este așa, se poate returna un array cu toate e-mailurile sau se poate prelucra primul, iar restul — la următoarea trecere. În general, totul depinde de sarcină.
Adăugăm la hook două funcții auxiliare: pentru descărcarea unui fișier și pentru descărcarea unui fișier dintr-un link din e-mail. Apropo, acestea pot fi extrase într-un operator, în funcție de frecvența utilizării acestei funcționalități. Ce altceva să mai adaug la hook, din nou, depinde de sarcină: dacă în e-mail vin fișiere, atunci putem descărca aplicațiile atașate, dacă datele vin în e-mail, atunci trebuie să parsăm e-mailul etc. În cazul meu, e-mailul vine cu un singur link către un arhivă, pe care trebuie să o plasăm într-un loc specific și să lansăm procesul de prelucrare ulterioară.
def download_from_url(self, url, path, chunk_size=128):
"""
Metodă pentru descărcarea unui fișier
:param url: Adresa de descărcare
:type url: string
:param path: Unde să plasăm fișierul
:type path: string
:param chunk_size: Câte bytes să scriem
: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):
"""
Metodă pentru descărcarea unui fișier dintr-un link din e-mail
:param mail_id: Identificatorul e-mailului
:type mail_id: string
:param path: Unde să plasăm fișierul
: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)Codul este simplu, așa că nu necesită explicații suplimentare. Voi menționa doar despre linia magică imap_conn_id. Apache Airflow stochează parametrii de conexiuni (username, parolă, adresă și alți parametri), la care se poate accesa printr-un identificator de tip string. Vizual, gestionarea conexiunilor arată astfel

Senzor pentru așteptarea datelor
Având în vedere că deja știm cum să ne conectăm și să obținem date din e-mail, acum putem scrie un senzor pentru a le aștepta. Nu am reușit să scriu simultan un operator care să proceseze datele, dacă acestea sunt disponibile, deoarece, pe baza datelor primite din e-mail, mai funcționează și alte procese, inclusiv cele care preiau date conexe din alte surse (API, telefonie, metrici web etc.). Voi da un exemplu. În sistemul CRM a apărut un nou utilizator și noi nu știm încă despre UUID-ul său. Atunci, la încercarea de a obține date de la telefonia SIP, vom primi apeluri legate de UUID-ul său, dar nu le putem salva și utiliza corect. În astfel de situații, este important să avem în vedere dependența datelor, mai ales dacă acestea provin din surse diferite. Acestea, desigur, sunt măsuri insuficiente pentru a menține integritatea datelor, dar, în unele cazuri, sunt necesare. Și, în plus, a ocupa resursele fără rost nu este rațional.
Astfel, senzorul nostru va lansa nodurile următoare ale grafului, dacă există informații proaspete în e-mail, precum și va marca informațiile anterioare ca fiind inadecvate.
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 TrueObținem și folosim datele
Pentru a obține și procesa datele, se poate scrie un operator separat sau se pot utiliza cele gata făcute. Întrucât logică este în prezent trivială — a prelua date dintr-un e-mail, pentru exemplu, sugerez un PythonOperator standard.
de la airflow.models import DAG
de la airflow.operators.python_operator import PythonOperator
de la airflow.sensors.my_plugin import MailSensor
de la my_plugin.hooks.imap_hook import IMAPHook
start_date = datetime(2020, 4, 4)
# Configurarea standard a graficului
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",
)
# Definirea senzorului
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",
)
# Funcția pentru obținerea datelor din email
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("Mailboxul este gol")
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)
# Descrierea celorlalte noduri ale graficului
...
# Stabilirea legăturii în grafic
mail_check_sensor >> prepare_mail_data
prepare_data >> ...
# Descrierea celorlalte fluxuri de controlApropo, dacă emailul dvs. corporativ este tot pe mail.ru, căutarea emailurilor după subiect, expeditor, etc. vă va fi inaccesibilă. Ei au promis asta în 2016, dar se pare că s-au răzgândit. Am rezolvat această problemă, creând un folder separat pentru emailurile necesare și configurând un filtru în interfața web pentru emailuri. Astfel, în acest folder ajung doar emailurile necesare, iar condițiile de căutare în cazul meu sunt simple (UNSEEN).
În rezumat, avem următoarea secvență: verificăm dacă există emailuri noi care îndeplinesc condițiile, iar dacă există, descărcăm arhiva prin linkul din ultimul email.
Sub ultimele puncte de suspensie este subînțeles că această arhivă va fi dezarhivată, datele din arhivă vor fi curățate și procesate, iar în final, toate aceste date vor fi trimise mai departe pe conveyorul procesului ETL, dar aceasta deja depășește subiectul articolului. Dacă a fost interesant și util, voi continua cu plăcere să descriu soluții ETL și părțile lor pentru Apache Airflow.
Sursa: habr.com
