
Egal wie stark sich die Technologien entwickeln, stets zieht eine Reihe veralteter Ansätze hinterher. Dies kann durch einen sanften Übergang, menschliche Faktoren, technologische Notwendigkeiten oder andere Dinge bedingt sein. Im Bereich der Datenverarbeitung sind die Datenquellen in dieser Hinsicht besonders aussagekräftig. So sehr wir uns danach sehnen, dem zu entkommen, solange ein Teil der Daten über Messenger und E-Mails versendet wird, ganz zu schweigen von den noch archaischeren Formaten. Ich lade ein, unter dem Cut einen der Ansätze für Apache Airflow zu betrachten, der veranschaulicht, wie man Daten aus E-Mails abrufen kann.
Vorgeschichte
Viele Daten werden nach wie vor über E-Mail übertragen, angefangen bei zwischenmenschlicher Kommunikation bis hin zu Standards für die Interaktion zwischen Unternehmen. Es ist von Vorteil, wenn man eine Schnittstelle für den Datenabruf schreiben kann oder Leute im Büro hat, die diese Informationen in bequemere Quellen eingeben, aber oft gibt es einfach nicht die Möglichkeit dafür. Die spezifische Aufgabe, mit der ich konfrontiert war, war die Anbindung eines bekannten CRM-Systems an ein Datenspeicher und anschließend an ein OLAP-System. Historisch gesehen war die Nutzung dieses Systems für unser Unternehmen in einem bestimmten Geschäftsbereich von Vorteil. Daher wünschten sich alle, mit Daten aus diesem externen System umgehen zu können. Zunächst wurde die Möglichkeit untersucht, Daten über die offene API abzurufen. Leider deckte die API nicht alle benötigten Daten ab, und einfach gesagt, war sie in vielerlei Hinsicht unzureichend, und der technische Support war nicht bereit oder in der Lage, den benötigten umfassenderen Funktionsumfang bereitzustellen. Stattdessen bot dieses System die Möglichkeit, fehlende Daten regelmäßig per E-Mail als Link zum Herunterladen des Archivs zu erhalten.
Es ist zu beachten, dass dies nicht der einzige Fall war, in dem das Unternehmen Daten aus E-Mails oder Messengern sammeln wollte. In diesem speziellen Fall konnten wir jedoch keinen Einfluss auf das externe Unternehmen nehmen, das einen Teil der Daten nur auf diese Weise bereitstellt.
Apache Airflow
Zur Erstellung von ETL-Prozessen verwenden wir am häufigsten Apache Airflow. Damit der Leser, der mit dieser Technologie nicht vertraut ist, besser versteht, wie dies im Kontext und insgesamt aussieht, werde ich ein paar einführende Informationen geben.
Apache Airflow ist eine Open-Source-Plattform, die zur Erstellung, Ausführung und Überwachung von ETL (Extract-Transform-Loading)-Prozessen in Python verwendet wird. Das Hauptkonzept in Airflow ist ein gerichteter azyklischer Graph, in dem die Knoten des Graphen konkrete Prozesse sind und die Kanten des Graphen den Fluss von Steuer- oder Informationsdaten darstellen. Ein Prozess kann einfach eine beliebige Python-Funktion aufrufen oder eine komplexere Logik in Form der sequenziellen Ausführung mehrerer Funktionen im Kontext einer Klasse haben. Für die häufigsten Operationen gibt es bereits viele vorgefertigte Lösungen, die als Prozesse verwendet werden können. Zu diesen Lösungen gehören:
- Operatoren – zum Verschieben von Daten von einem Ort an einen anderen, zum Beispiel von einer Datenbanktabelle in ein Data Warehouse;
- Sensoren – um auf das Eintreten eines bestimmten Ereignisses zu warten und den Fluss der Steuerung zu den nachfolgenden Knoten des Graphen zu leiten;
- Hooks – für niedrigere Operationen, zum Beispiel zum Abrufen von Daten aus einer Datenbanktabelle (werden in Operatoren verwendet);
- usw.
Es wäre nicht sinnvoll, Apache Airflow in diesem Artikel im Detail zu beschreiben. Kurze Einführungen können angesehen werden oder .
Hook zum Abrufen von Daten
Zuerst müssen wir einen Hook schreiben, mit dem wir:
- uns mit dem E-Mail-Server verbinden;
- die benötigte E-Mail finden;
- Daten aus der E-Mail abrufen.
from airflow.hooks.base_hook import BaseHook
import imaplib
import logging
class IMAPHook(BaseHook):
def __init__(self, imap_conn_id):
"""
IMAP-Hook zum Abrufen von Daten von E-Mails
:param imap_conn_id: E-Mail-Verbindungs-ID
:type imap_conn_id: string
"""
self.connection = self.get_connection(imap_conn_id)
self.mail = None
def authenticate(self):
"""
Verbindung zur E-Mail herstellen
"""
mail = imaplib.IMAP4_SSL(self.connection.host)
response, detail = mail.login(user=self.connection.login, password=self.connection.password)
if response != "OK":
raise AirflowException("Anmeldung fehlgeschlagen")
else:
self.mail = mail
def get_last_mail(self, check_seen=True, box="INBOX", condition="(UNSEEN)"):
"""
Methode zum Abrufen der ID der letzten E-Mail,
die den Suchkriterien entspricht
:param check_seen: Letzte E-Mail als gelesen markieren
:type check_seen: bool
:param box: Bezeichnung des Postfachs
:type box: string
:param condition: Suchbedingungen für 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("Folgende E-Mails wurden im Postfach gefunden: " + str(mail_ids))
if not mail_ids:
logging.info("Keine neuen E-Mails gefunden")
return None
mail_id = mail_ids[0]
# wenn es mehrere solcher E-Mails gibt
if len(mail_ids) > 1:
# die übrigen als gelesen markieren
for id in mail_ids:
self.mail.store(id, "+FLAGS", "\Seen")
# die letzte zurückgeben
mail_id = mail_ids[-1]
# Soll die letzte E-Mail als gelesen markiert werden?
if not check_seen:
self.mail.store(mail_id, "-FLAGS", "\Seen")
return mail_idDie Logik ist folgende: Wir verbinden uns, finden die letzte aktuelle E-Mail, und wenn es andere gibt, ignorieren wir diese. Diese Funktion wird genau aus diesem Grund verwendet, da spätere E-Mails alle Daten der früheren enthalten. Ist das nicht der Fall, könnte man ein Array aller E-Mails zurückgeben oder die erste bearbeiten und die anderen beim nächsten Durchlauf. Im Grunde hängt alles wie gewohnt von der Aufgabe ab.
Wir fügen dem Hook zwei Hilfsfunktionen hinzu: eine zum Herunterladen einer Datei und eine zum Herunterladen einer Datei über einen Link in einer E-Mail. Übrigens können sie in einen Operator ausgelagert werden, dies hängt von der Häufigkeit der Nutzung dieser Funktionalität ab. Was sonst noch in den Hook geschrieben werden sollte, hängt wiederum von der Aufgabe ab: Wenn in der E-Mail sofort Dateien enthalten sind, kann man die Anhänge herunterladen; wenn Daten in der E-Mail ankommen, muss die E-Mail geparsed werden usw. In meinem Fall kommt die E-Mail mit einem Link zu einem Archiv, das ich an einem bestimmten Ort ablegen und den weiteren Verarbeitungsprozess starten muss.
def download_from_url(self, url, path, chunk_size=128):
"""
Methode zum Herunterladen einer Datei
:param url: Die Download-Adresse
:type url: string
:param path: Wo die Datei gespeichert werden soll
:type path: string
:param chunk_size: Wie viele Bytes auf einmal geschrieben werden
: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):
"""
Methode zum Herunterladen einer Datei über einen Link in der E-Mail
:param mail_id: Die ID der E-Mail
:type mail_id: string
:param path: Wo die Datei gespeichert werden soll
: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)Der Code ist einfach, daher benötigt er wahrscheinlich keine zusätzlichen Erklärungen. Ich werde nur über die magische Zeile imap_conn_id sprechen. Apache Airflow speichert Verbindungseinstellungen (Benutzername, Passwort, Adresse und andere Parameter), auf die anhand einer String-ID zugegriffen werden kann. Visuell sieht die Verwaltung von Verbindungen so aus

Sensor zum Warten auf Daten
Da wir bereits wissen, wie man sich verbindet und Daten aus der E-Mail abruft, können wir jetzt einen Sensor schreiben, um auf diese zu warten. Es gelang mir nicht, sofort einen Operator zu schreiben, der die Daten verarbeitet, wenn sie vorhanden sind, da auf der Grundlage der empfangenen Daten aus der E-Mail auch andere Prozesse laufen, die verwandte Daten aus anderen Quellen abrufen (API, Telefonie, Webmetriken usw.). Ich gebe ein Beispiel. In der CRM-System ist ein neuer Nutzer erschienen, und wir wissen noch nichts über seine UUID. Wenn wir also versuchen, Daten von der SIP-Telefonie abzurufen, erhalten wir Anrufe, die an seine UUID gebunden sind, können sie aber nicht korrekt speichern und verwenden. In solchen Fragen ist es wichtig, die Abhängigkeit der Daten im Auge zu behalten, insbesondere wenn sie aus verschiedenen Quellen stammen. Das sind zwar unzureichende Maßnahmen zur Wahrung der Datenintegrität, aber in manchen Fällen notwendig. Es ist auch nicht sinnvoll, Ressourcen nutzlos zu besetzen.
Unser Sensor wird die nachfolgenden Knoten des Grafen auslösen, wenn es frische Informationen in der E-Mail gibt, und gleichzeitig die vorherigen Informationen als nicht mehr aktuell kennzeichnen.
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="Posteingang", 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 TrueDaten abrufen und verwenden
Um Daten zu erhalten und zu verarbeiten, kann man einen separaten Operator schreiben oder bestehende verwenden. Da die Logik momentan trivial ist – die Daten aus der E-Mail abzurufen, schlage ich als Beispiel den Standard-PythonOperator vor.
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)
# Standardkonfiguration des Graphen
args = {
"owner": "example",
"start_date": start_date,
"email": ["home@home.de"],
"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 definieren
mail_check_sensor = MailSensor(
task_id="check_new_emails",
poke_interval=10,
conn_id="mail_conn_id",
timeout=10,
soft_fail=True,
box="mein_postfach",
dag=dag,
mode="poke",
)
# Funktion zur Datenbeschaffung aus der Mail
def prepare_mail():
imap_hook = IMAPHook("mail_conn_id")
mail_id = imap_hook.get_last_mail(check_seen=True, box="mein_postfach")
if mail_id is None:
raise AirflowException("Leeres Postfach")
conn.download_mail_href_attachment(mail_id, ".\/pfad.zip")
prepare_mail_data = PythonOperator(task_id="prepare_mail_data", default_args=args, dag=dag, python_callable=prepare_mail)
# Beschreibung der anderen Knoten des Graphen
...
# Verknüpfung im Graphen festlegen
mail_check_sensor >> prepare_mail_data
prepare_data >> ...
# Beschreibung der anderen SteuerflüsseÜbrigens, wenn Ihr Firmen-E-Mail-Account ebenfalls bei mail.ru ist, wird Ihnen die Suche nach E-Mails nach Betreff, Absender usw. nicht zur Verfügung stehen. Sie hatten bereits 2016 versprochen, dies einzuführen, aber anscheinend haben sie es sich anders überlegt. Ich habe dieses Problem gelöst, indem ich ein separates Verzeichnis für die benötigten E-Mails erstellt und im Web-Interface meiner E-Mail einen Filter für die gewünschten E-Mails eingerichtet habe. So gelangen nur die benötigten E-Mails in diesen Ordner und die Suchbedingungen sind in meinem Fall einfach (UNSEEN).
Zusammenfassend haben wir die folgende Abfolge: Wir prüfen, ob neue E-Mails vorliegen, die den Bedingungen entsprechen, wenn ja, laden wir das Archiv über den Link aus der letzten E-Mail herunter.
Unter den letzten Auslassungen ist zu erwähnen, dass dieses Archiv entpackt, die Daten aus dem Archiv gereinigt und bearbeitet werden und schließlich alles weiter in den ETL-Prozess fließt, aber das geht bereits über den Rahmen des Artikels hinaus. Wenn es interessant und nützlich war, werde ich gerne weiterhin ETL-Lösungen und deren Teile für Apache Airflow beschreiben.
Quelle: habr.com
