
Nonostante l'evoluzione delle tecnologie, un gran numero di approcci obsoleti continua a persistere. Questo può essere dovuto a una transizione graduale, a fattori umani, necessità tecnologiche o altri motivi. Nel campo dell'elaborazione dei dati, le fonti di dati sono particolarmente rappresentative di questo fenomeno. Per quanto possiamo desiderare di eliminarle, fino a quando una parte dei dati viene inviato tramite messaggistica e email, senza contare i formati più arcaici, questo problema persisterà. Vi invito a scoprire di seguito una delle soluzioni per Apache Airflow che illustra come è possibile prelevare dati dalle email.
Contesto
Molti dati vengono ancora trasmessi via email, dalle comunicazioni interpersonali agli standard di interazione tra aziende. È utile se si riesce a scrivere un'interfaccia per ottenere i dati o a coinvolgere persone in ufficio per inserire queste informazioni in fonti più comode, ma spesso tale opportunità può semplicemente non esserci. Il compito specifico con cui mi sono trovato a confrontarmi è stato collegare un noto sistema CRM a un repository di dati e, successivamente, a un sistema OLAP. Storicamente, per la nostra azienda, l'utilizzo di questo sistema era comodo in un'area specifica del business. Pertanto, era molto desiderato poter operare con i dati anche da questo sistema esterno. Prima di tutto, è stata esaminata la possibilità di ottenere dati attraverso un API aperto. Sfortunatamente, l'API non copriva l'acquisizione di tutti i dati necessari e, per dirla in modo semplice, presentava molte anomalie, e il supporto tecnico non ha voluto o potuto essere d'aiuto per fornire funzionalità più esaustive. Tuttavia, questo sistema forniva la possibilità di ricevere periodicamente i dati mancanti via email sotto forma di un collegamento per il download dell'archivio.
È importante notare che questo non era l'unico caso in cui l'azienda voleva raccogliere dati dalle e-mail o dai messenger. Tuttavia, in questo caso non potevamo influenzare la società terza che fornisce parte dei dati solo in questo modo.
Apache Airflow
Per costruire i processi ETL, utilizziamo più frequentemente Apache Airflow. Affinché il lettore, non familiare con questa tecnologia, possa comprendere meglio come si presenta nel contesto e in generale, descriverò un paio di aspetti introduttivi.
Apache Airflow è una piattaforma open source utilizzata per costruire, eseguire e monitorare processi ETL (Extract-Transform-Loading) in Python. Il concetto principale in Airflow è il grafo aciclico orientato, dove i nodi del grafo rappresentano processi specifici e i bordi del grafo rappresentano il flusso di controllo o informazioni. Un processo può semplicemente chiamare qualsiasi funzione Python oppure avere una logica più complessa composta dalla chiamata sequenziale di più funzioni nel contesto di una classe. Per le operazioni più comuni, ci sono già molti componenti pronti che possono essere utilizzati come processi. Tra questi ci sono:
- operatori — per il trasferimento di dati da un luogo all'altro, ad esempio da una tabella del DB a un data warehouse;
- sensori — per attendere il verificarsi di un determinato evento e dirigere il flusso di controllo verso i successivi nodi del grafo;
- hook — per operazioni a un livello inferiore, ad esempio per recuperare dati da una tabella del DB (utilizzati negli operatori);
- ecc.
Non sarebbe utile descrivere Apache Airflow in dettaglio in questo articolo. È possibile visualizzare brevi introduzioni o .
Hook per il recupero dei dati
In primo luogo, per risolvere il problema dobbiamo scrivere un hook che ci permetta di:
- collegarci alla posta elettronica;
- trovare l'email desiderata;
- recuperare i dati dall'email.
from airflow.hooks.base_hook import BaseHook
import imaplib
import logging
class IMAPHook(BaseHook):
def __init__(self, imap_conn_id):
"""
Hook IMAP per ottenere dati dalle email
:param imap_conn_id: Identificativo della connessione email
:type imap_conn_id: string
"""
self.connection = self.get_connection(imap_conn_id)
self.mail = None
def authenticate(self):
"""
Ci connettiamo all'email
"""
mail = imaplib.IMAP4_SSL(self.connection.host)
response, detail = mail.login(user=self.connection.login, password=self.connection.password)
if response != "OK":
raise AirflowException("Accesso non riuscito")
else:
self.mail = mail
def get_last_mail(self, check_seen=True, box="INBOX", condition="(UNSEEN)"):
"""
Metodo per ottenere l'identificativo dell'ultima email,
che soddisfa le condizioni di ricerca
:param check_seen: Segnare l'ultima email come letta
:type check_seen: bool
:param box: Nome della casella
:type box: string
:param condition: Condizioni di ricerca delle email
:type condition: string
"""
self.authenticate()
self.mail.select(mailbox=box)
response, data = self.mail.search(None, condition)
mail_ids = data[0].split()
logging.info("Nella casella sono state trovate le seguenti email: " + str(mail_ids))
if not mail_ids:
logging.info("Nessuna nuova email trovata")
return None
mail_id = mail_ids[0]
# se ci sono più email
if len(mail_ids) > 1:
# segna le altre come lette
for id in mail_ids:
self.mail.store(id, "+FLAGS", "\Seen")
# restituisce l'ultima
mail_id = mail_ids[-1]
# bisogna segnare l'ultima come letta?
if not check_seen:
self.mail.store(mail_id, "-FLAGS", "\Seen")
return mail_idLa logica è la seguente: ci connettiamo, troviamo l'ultima email più aggiornata, e se ce ne sono altre le ignoriamo. Questo metodo è usato perché le email più recenti contengono tutti i dati delle precedenti. Se non è così, possiamo restituire un array di tutte le email o elaborare la prima e le altre alla prossima iterazione. In generale, tutto dipende dall'obiettivo.
Aggiungiamo due funzioni ausiliarie all'hook: una per il download del file e l'altra per il download del file tramite il link presente nell'email. A proposito, possono essere estratte in un operatore, a seconda della frequenza d'uso di questa funzionalità. Cosa aggiungere ulteriormente all'hook dipende anche dall'obiettivo: se nell'email ci sono file allegati, possiamo scaricarli; se i dati arrivano nell'email, dobbiamo analizzarla, ecc. Nel mio caso, l'email arriva con un link a un archivio che devo posizionare in un luogo specifico e avviare il processo di elaborazione.
def download_from_url(self, url, path, chunk_size=128):
"""
Metodo per scaricare un file
:param url: Indirizzo di download
:type url: string
:param path: Dove salvare il file
:type path: string
:param chunk_size: Byte per scrittura
: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):
"""
Metodo per scaricare un file dal link in una mail
:param mail_id: Identificativo della mail
:type mail_id: string
:param path: Dove salvare il 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)Il codice è semplice e quindi non richiede ulteriori spiegazioni. Dirò solo qualcosa sulla magica stringa imap_conn_id. Apache Airflow memorizza i parametri di connessione (login, password, indirizzo e altri parametri), ai quali si può accedere tramite un identificatore stringa. La gestione delle connessioni appare così.

Sensore per l'attesa dei dati
Poiché sappiamo già come collegarci e ricevere dati dalle e-mail, ora possiamo scrivere un sensore per attenderli. Non sono riuscito a scrivere subito un operatore che elaborasse i dati se disponibili, poiché in base ai dati ricevuti dalle e-mail iterano altri processi, compresi quelli che estraggono dati correlati da altre fonti (API, telefonia, web analytics, ecc.). Faccio un esempio. Se in un sistema CRM compare un nuovo utente e non conosciamo ancora il suo UUID, quando tentiamo di ottenere i dati dalla telefonia SIP, riceveremo chiamate collegate al suo UUID, ma non saremo in grado di conservarle e utilizzarle correttamente. In queste situazioni è fondamentale tenere presente la dipendenza dei dati, soprattutto se provengono da fonti diverse. Queste, ovviamente, sono misure insufficienti per preservare l'integrità dei dati, ma in alcuni casi sono necessarie. Inoltre, occupare risorse inutilmente non è razionale.
Pertanto, il nostro sensore attiverà i nodi successivi del grafo se ci sono nuove informazioni nelle e-mail, e contrassegnerà come obsolete le informazioni precedenti.
da 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 TrueRiceviamo e utilizziamo i dati
Per ottenere e elaborare i dati, puoi scrivere un operatore separato o utilizzare quelli predefiniti. Poiché la logica è attualmente semplice — estrarre i dati da un'email — per questo esempio ti propongo di utilizzare il PythonOperator standard.
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)
# Configurazione standard del grafo
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",
)
# Definizione del sensore
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",
)
# Funzione per estrarre i dati dall'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("Mailbox vuota")
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)
# Descrizione delle altre vertici del grafo
...
# Impostiamo i collegamenti nel grafo
mail_check_sensor >> prepare_mail_data
prepare_data >> ...
# Descrizione degli altri flussi di controlloTra l'altro, se la tua email aziendale è anch'essa su mail.ru, non potrai cercare le email per oggetto, mittente, ecc. Promisero di introdurre questa funzione nel lontano 2016, ma evidentemente hanno cambiato idea. Ho risolto il problema creando una cartella separata per le email necessarie e impostando un filtro nell'interfaccia web della posta per le email desiderate. In questo modo, in questa cartella finiscono solo le email necessarie e i criteri per cercare nel mio caso sono semplicemente (UNSEEN).
In sintesi, abbiamo la seguente sequenza: verifichiamo se ci sono nuove email che soddisfano i criteri; se ci sono, scarichiamo l'archivio dal link dell'ultima email.
Sotto i punti di sospensione finali viene omesso che quest'archivio sarà estratto, i dati dell'archivio saranno puliti e elaborati, e alla fine tutto questo andrà oltre nella catena del processo ETL, ma questo già esula dall'argomento dell'articolo. Se è stato interessante e utile, sarò felice di continuare a descrivere soluzioni ETL e le loro parti per Apache Airflow.
Fonte: habr.com
