
Kuigi tehnoloogia areneb kiiresti, järgneb sellele alati rida vananenud lähenemisi. Selle põhjuseks võivad olla sujuvad üleminekud, inimfaktor, tehnoloogilised vajadused või midagi muud. Andmete töötlemise valdkonnas on kõige silmapaistvamad andmeallikad. Ükskõik kui palju me unistame sellest vabaneda, saadetakse osa andmeid endiselt sõnumirakendustes ja e-kirjades, rääkimata veelgi archailisematest formaatidest. Kutsub üles allolevat arutama üht võimalust Apache Airflow jaoks, mis illustreerib, kuidas andmeid e-kirjadest välja saada.
Eelalugu
Paljud andmed edastatakse endiselt e-posti kaudu, alates omavahelistest suhtlustest kuni ettevõtetevaheliste suhtluse standarditeni. Hea, kui on võimalik andmete saamiseks kirjutada liides või koosolekuteks inimesi kontorisse paigutada, kes need andmed mugavatesse allikatesse sisestavad, kuid sageli selliseks võimaluseks lihtsalt ei ole. Konkreetne ülesanne, millega mina silmitsi seisma pidin, oli tuntud CRM-süsteemi ühendamine andmehoidla ning seejärel OLAP-süsteemiga. Ajalooliselt on meie ettevõtte jaoks olnud selle süsteemi kasutamine mugav teatud ärivaldkonnas. Seetõttu soovisid kõik väga osata manipuleerida andmetega ka sellest kolmandast süsteemist. Esiteks uuriti loomulikult andmete saamise võimalust avatud API kaudu. Kahjuks ei katnud API kõigi vajalike andmete saamist ning, lihtsas keeles öeldes, oli see paljuski vigane, ning tehniline tugi ei soovinud ega suutnud pakkuda rohkem ammendavat funktsionaalsust. Kuid see süsteem pakkus võimalust perioodiliselt saada puudulikud andmed e-posti teel lingina arhiivi laadimiseks.
On oluline märkida, et see ei olnud ainus juhtum, mille puhul ettevõte soovis koguda andmeid e-kirjadest või sõnumiteenusest. Kuid antud juhul ei saanud me mõjutada kolmandat ettevõtet, mis pakub osa andmeid ainult sedaviisi.
Apache Airflow
ETL-protsesside ehitamiseks kasutame kõige sagedamini Apache Airflow't. Selleks, et lugeja, kes pole selle tehnoloogiaga tuttav, paremini mõistaks, kuidas see kontekstis välja näeb, kirjeldan paar sissejuhatavat punkt.
Apache Airflow on avatud lähtekoodiga platvorm, mida kasutatakse ETL (Extract-Transform-Loading) protsesside ehitamiseks, käitamiseks ja jälgimiseks Pythonis. Peamine mõiste Airflow's on suunatud tsükliline graaf, kus graafi tipud on konkreetsed protsessid ja graafi servad on juhtimis- või infovoog. Protsess võib lihtsalt kutsuda esile ükskõik millise Python-funktsiooni või olla keerukama loogikaga, mis koosneb mitme funktsiooni järjestikusest kutsumisest klassi kontekstis. Kõige sagedasemateks toiminguteks on juba palju valmislahendusi, mida saab kasutada protsessidena. Nende hulka kuuluvad:
- operaatorid – andmete edastamiseks ühest kohast teise, näiteks andmebaasi tabelist andmete ladustamisse;
- sensorid – teatud sündmuse toimumise ootamiseks ja juhtimisvoo suunamiseks graafi järgmistele tippudele;
- häkid – madalama taseme operatsioonide jaoks, näiteks andmete saamiseks andmebaasi tabelist (kasutatakse operaatorites);
- Kui kaua aega kulub väljastamiseks?
Apache Airflow'i detailselt kirjeldamine selles artiklis ei ole eesmärgipärane. Lühikesi sissejuhatusi saab vaadata või .
Andmete saamise häkk
Esiteks, ülesande lahendamiseks tuleb kirjutada häkk, mille abil me saaksime:
- ühenduda e-postiga;
- leida vajalik kiri;
- saada andmed kirjast.
from airflow.hooks.base_hook import BaseHook
import imaplib
import logging
class IMAPHook(BaseHook):
def __init__(self, imap_conn_id):
"""
IMAP hook e-postiandmete saamiseks
: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 e-posti
"""
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 meilisõnumi tuvastamiseks,
mis vastab otsingutingimustele
:param check_seen: Märkida viimane meilisõnum loetuks
:type check_seen: bool
:param box: Kasti nimetus
:type box: string
:param condition: Meilisõnumite 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("Kastis leiti järgmised meilisõnumid: " + str(mail_ids))
if not mail_ids:
logging.info("Uusi meilisõnumeid ei leitud")
return None
mail_id = mail_ids[0]
# kui selliseid meilisõnumeid on mitu
if len(mail_ids) > 1:
# märgime ülejäänud loetuks
for id in mail_ids:
self.mail.store(id, "+FLAGS", "\Seen")
# tagastame viimase
mail_id = mail_ids[-1]
# kas tuleb viimane märkida loetuks
if not check_seen:
self.mail.store(mail_id, "-FLAGS", "\Seen")
return mail_idLoogika on järgmine: ühendume, leiame kõige värskema kirja, kui on teisi - ignoreerime neid. Kasutatakse just sellist funktsiooni, kuna hilisemad kirjateated sisaldavad kõiki varasemate andmeid. Kui see ei kehti, siis võib tagastada kõigi kirjade massiivi või töödelda esimest, ning teised järgmise läbimise korral. Üldiselt sõltub kõik ülesandest.
Lisame hook'ile kaks abifunktsiooni: faili allalaadimiseks ja faili allalaadimiseks lingilt kirjast. Pean ütlema, et need saab välja tuua operaatorisse, sõltuvalt sellest, kui sageli seda funktsionaalsust kasutatakse. Mida veel hook'ile lisada, sõltub jälle ülesandest: kui kirjas on kohe failid, siis saab allalaadida kirja manuseid, kui andmed tulevad kirjas, tuleb kirja analüüsida jne. Minu puhul tuleb kiri ühe arhiivi lingiga, mille ma pean panema kindlasse kohta ja käivitama edasise töötlemise.
def download_from_url(self, url, path, chunk_size=128):
"""
Meetod failide allalaadimiseks
:param url: Laadimisadresse
:type url: string
:param path: Kuhu fail panna
:type path: string
:param chunk_size: Mitu baidi 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):
"""
Meetod failide allalaadimiseks lingilt kirjast
:param mail_id: Kirja identifikaator
:type mail_id: string
:param path: Kuhu fail 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 ei vaja see tõenäoliselt täiendavaid selgitusi. Räägin vaid maagilisest reast imap_conn_id. Apache Airflow salvestab ühendusparameetrid (kasutajanimi, parool, aadress ja muud parameetrid), millele saab juurde pääseda stringi identifikaatori kaudu. Visuaalselt näeb ühenduste haldamine välja järgmiselt

Andmete ootamise sensor
Kuna me oskame juba e-kirjadega ühendust võtta ja andmeid saada, saame nüüd kirjutada sensori, mis neid ootab. Kohe kirjutada operaator, mis töötleb andmeid, kui need on olemas, ei õnnestunud, kuna e-kirjast saadud andmete põhjal töötavad ka teised protsessid, sealhulgas need, mis võtavad seotud andmeid teistest allikatest (API, telefoni- ja veebianalüütika jne). Olgu üks näide. Kui CRM-süsteemis ilmub uus kasutaja, ei tea me veel tema UUID-d. Siis, kui proovime SIP-telefonilt andmeid saada, saame kõnesid, mis on seotud tema UUID-ga, kuid ei suuda neid õigesti salvestada ja kasutada. Selles osas on oluline arvestada andmete sõltuvust, eriti kui need on pärit erinevatest allikatest. See ei ole loomulikult piisav meede andmete terviklikkuse säilitamiseks, kuid teatud juhtudel on see vajalik. Ja ressursse ka tühja koormata ei ole ratsionaalne.
Nii hakkab meie sensor käivitama järgmisi graafi tippe, kui postkastis on värske teave, ja märkima eelneva teabe aegunuks.
import airflow.sensors.base_sensor_operator as 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 TrueSaame ja kasutame andmeid
Andmete saamiseks ja töötlemiseks saab kirjutada eraldi operaatori või kasutada olemasolevaid. Kuna loogika on hetkel triviaalne — andmete saamine kirjast, pakun näiteks tavalist PythonOperatorit
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)
# Graafi standardkonfiguratsioon
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 sensori
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, et 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)
# Üksuste osade kirjeldus
...
# Määrame seose graafikus
mail_check_sensor >> prepare_mail_data
prepare_data >> ...
# Üksuste osade kirjeldusMuide, kui teie ettevõtte e-post on samuti mail.ru, siis ei ole teil võimalik otsida kirju teema, saatja jne järgi. Nad lubasid seda juba 2016. aastal, kuid ilmselt on nad meelt muutnud. Lahendasin selle probleemi, loobudes eraldi kausta vajalike kirjade jaoks ja seadistades veebiliideses e-posti jaoks filtrid sobivatele kirjadele. Nii pääsevad sellesse kausta ainult vajalikud kirjad ja minu otsingutingimused on lihtsalt (UNSEEN).
Kokkuvõtteks, meil on järgmine järjestus: kontrollime, kas on uusi kirju, mis vastavad tingimustele, kui on, siis laadime viimase kirja lingilt arhiivi alla.
Viimaste kolme punkti taga jääb mainimata, et see arhiiv avatakse, arhiivi andmed puhastatakse ja töödeldakse ning lõpuks saadetakse see edasi ETL protsessi konveierile, kuid see ületab juba artikli teema piire. Kui see tundus huvitav ja kasulik, siis kirjutan meeleldi edasi ETL lahendustest ja nende osadest Apache Airflow jaoks.
Allikas: habr.com
