
Sido zhvillimet teknologjike, gjithmonë pas tyre qëndrojnë qasje të vjetra. Kjo mund të jetë e lidhur me një kalim të ngadaltë, faktorë njerëzorë, nevojat teknologjike, ose diçka tjetër. Në fushën e përpunimit të të dhënave, burimet e të dhënave janë më të tregueshme në këtë aspekt. Sa do që të dëshirojmë të shmangim këtë, një pjesë e të dhënave ende dërgohet përmes mesazheve dhe email-eve, pa përmendur format më arkaike. Ju ftoj të eksploroni një nga opsionet për Apache Airflow, që ilustron si mund të tërheqim të dhëna nga email-et.
Pas historia
ShumĂ« tĂ« dhĂ«na akoma transferohen pĂ«rmes postĂ«s elektronike, duke filluar nga komunikimet ndĂ«rpersonale dhe pĂ«rfunduar me standardet e ndĂ«rveprimit midis kompanive. ĂshtĂ« mirĂ« nĂ«se mund tĂ« krijosh njĂ« ndĂ«rfaqe pĂ«r marrjen e tĂ« dhĂ«nave ose tĂ« vendosĂ«sh njerĂ«z nĂ« zyrĂ« qĂ« do tĂ« fusin kĂ«tĂ« informacion nĂ« burime mĂ« tĂ« pĂ«rshtatshme, por shpesh kjo mundĂ«si nuk mund tĂ« jetĂ« e disponueshme. Detyra specifike me tĂ« cilĂ«n u pĂ«rballa Ă«shtĂ« lidhu me njĂ« sistem CRM tĂ« njohur me njĂ« depo tĂ« dhĂ«nash, dhe mĂ« pas me njĂ« sistem OLAP. Historikisht, pĂ«rdorimi i kĂ«tij sistemi ishte i pĂ«rshtatshĂ«m pĂ«r njĂ« fushĂ« tĂ« veçantĂ« tĂ« biznesit nĂ« kompaninĂ« tonĂ«. Prandaj, tĂ« gjithĂ« dĂ«shironim shumĂ« tĂ« mund tĂ« operonim me tĂ« dhĂ«nat edhe nga ky sistem i jashtĂ«m. NĂ« radhĂ« tĂ« parĂ«, sigurisht, u studiuar mundĂ«sia e marrjes sĂ« tĂ« dhĂ«nave nga njĂ« API tĂ« hapur. FatkeqĂ«sisht, API nuk mbulonte marrjen e tĂ« dhĂ«nave tĂ« nevojshme, dhe, duke e shprehur nĂ« mĂ«nyrĂ« tĂ« thjeshtĂ«, ishte shumĂ« e pafavorshme, dhe mbĂ«shtetja teknike nuk dĂ«shiroi ose nuk arriti tĂ« ofronte njĂ« funksionalitet mĂ« tĂ« plotĂ«. Por, ky sistem ofronte mundĂ«sinĂ« e marrjes periodike tĂ« tĂ« dhĂ«nave tĂ« munguara nĂ«pĂ«rmjet postĂ«s qĂ« shpĂ«rndante njĂ« lidhje pĂ«r shkarkimin e arkivave.
Duhet theksuar se ky nuk ishte rasti i vetëm për të cilin biznesi dëshiron të mbledhë të dhëna nga email-at ose mesazhet. Megjithatë, në këtë rast ne nuk mund të ndikojmë në kompaninë e tretë, e cila ofron një pjesë të të dhënave 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 e panjohur me këtë teknologji të kuptojë më mirë se si duket në kontekst dhe në përgjithësi, do të përshkruaj disa hyrje.
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ë grafiku orientuar aklirik, ku majat e grafikëve janë procese të caktuara, ndërsa skajet e grafikëve janë rrjedha e kontrollit ose informacionit. Një proces mund të thërrasë thjesht një funksion Python, ose mund të ketë një logjikë më të komplikuar nga thirrja e rendit të disa funksioneve në kuadër të një klase. Për operacionet më të zakonshme, tashmë ka shumë përgatitje që mund të përdoren si procese. Këto përgatitje përfshijnë:
- operatorë - për transferimin e të dhënave nga një vend në një tjetër, për shembull nga tabela e DB në depo të dhënash;
- sensorë - për pritjen e një ngjarjeje të caktuar dhe për drejtimin e rrjedhës së kontrollit në piketat e mëpasshëm të grafit;
- hook - për operacione më të ulta, për shembull, për të marrë të dhëna nga tabela e DB (përdoren në operatorë);
- etj.
Të njihet Apache Airflow në detaje në këtë artikull nuk do të ishte e arsyeshme. Mund të shihni hyrje të shkurtra ose .
Hook për marrjen e të dhënave
Së pari, për të zgjidhur problemin, ne duhet të shkruajmë një hook me të cilin mund të:
- lidhemi me postën elektronike;
- të gjejmë mesazhin e nevojshëm;
- të marrim të dhëna nga mesazhi.
from 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 e-mail
:param imap_conn_id: Identifikuesi i lidhjes me e-mail
:type imap_conn_id: string
"""
self.connection = self.get_connection(imap_conn_id)
self.mail = None
def authenticate(self):
"""
Lidhemi me 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("Hyrja dështoi")
else:
self.mail = mail
def get_last_mail(self, check_seen=True, box="INBOX", condition="(UNSEEN)"):
"""
Metodë për të marrë identifikuesin e e-mailit të fundit,
që përmbush kushtet e kërkimit
:param check_seen: Të shënosh e-mailin e fundit si të lexuar
:type check_seen: bool
:param box: Emri i kuti
:type box: string
:param condition: Kushtet e kërkimit të e-maileve
: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 këto e-maila: " + str(mail_ids))
if not mail_ids:
logging.info("Nuk u gjetën e-maila të rinj")
return None
mail_id = mail_ids[0]
# nëse ka shumë e-maile të tillë
if len(mail_ids) > 1:
# shënojmë të tjerët si të lexuar
for id in mail_ids:
self.mail.store(id, "+FLAGS", "\Seen")
# kthejmë të fundit
mail_id = mail_ids[-1]
# a duhet të shënojmë të fundit si të lexuar
if not check_seen:
self.mail.store(mail_id, "-FLAGS", "\Seen")
return mail_idLogjika është e tillë: lidhemi, gjejmë emailin më të fundit dhe më aktual, nëse ka të tjerë - i injorojmë. Përdoret pikërisht kjo funksion, sepse emailët më të vonshëm përmbajnë të gjithë të dhënat e mëparshme. Nëse kjo nuk është e vërtetë, mund të kthejmë një masiv të të gjithë emaileve ose të përpunojmë të parin dhe të tjerët në kalimin e ardhshëm. Në përgjithësi, gjithçka varet nga detyra.
ShtojmĂ« dy funksione ndihmĂ«se nĂ« hook: pĂ«r shkarkimin e skedarĂ«ve dhe pĂ«r shkarkimin e skedarĂ«ve nga lidhja nĂ« email. PĂ«r ta thĂ«nĂ« ndryshe, ato mund tĂ« nxirren nĂ« operator, kjo varet nga frekuenca e pĂ«rdorimit tĂ« kĂ«tij funksionaliteti. ĂfarĂ« tjetĂ«r duhet shtuar nĂ« hook, pĂ«rsĂ«ri varet nga detyra: nĂ«se nĂ« email vijnĂ« menjĂ«herĂ« skedarĂ«, mund tĂ« shkarkojmĂ« aplikimet nĂ« email; nĂ«se tĂ« dhĂ«nat vijnĂ« nĂ« email, atĂ«herĂ« duhet tĂ« analizojmĂ« emailin dhe etj. NĂ« rastin tim, emaili vjen me njĂ« lidhje pĂ«r njĂ« arkiv, qĂ« 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):
"""
Metod për shkarkimin e një skedari
:param url: Adresa e shkarkimit
:type url: string
:param path: Mënyra për ta ruajtur skedarin
:type path: string
:param chunk_size: Sa 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):
"""
Metod për shkarkimin e një skedari nga linku në letër
:param mail_id: Identifikuesi i letrës
:type mail_id: string
:param path: Mënyra për ta ruajtur skedarin
: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ë, prandaj ndoshta nuk ka nevojë për shpjegime të tjera. Do të flas vetëm për vijën magjike imap_conn_id. Apache Airflow ruan parametrat e lidhjeve (emrin e përdoruesit, fjalëkalimin, adresën dhe parametrat e tjerë), të cilat mund të aksesohen përmes identifikuesit të vargut. Vizualisht, menaxhimi i lidhjeve duket kështu

Sensor për pritjen e të dhënave
Tani që ne tashmë dimë si të lidhemi dhe të marrim të dhëna nga posta, tani mund të shkruajmë një sensor për t'i pritur ato. Të shkruajmë menjëherë një operator që do të përpunojë të dhënat, kur ato janë të pranishme, nuk arrita ta bëj, pasi proceset e tjera po funksionojnë në bazë të të dhënave të marra nga posta, përfshirë ato që marrin të dhëna të lidhura nga burime të tjera (API, telefonia, analiza e webit, etj.). Do të jap një shembull. Në sistemin CRM u shfaq një përdorues i ri, dhe ne ende nuk e dimë në lidhje me UUID-në e tij. Atëherë, kur të përpiqemi të marrim të dhëna nga telefonia SIP, do të marrim thirrjet të lidhura me UUID-në e tij, por nuk do të jemi në gjendje t'i ruajmë dhe t'i përdorim ato në mënyrë të saktë. Në këto raste, është e rëndësishme të kemi parasysh varësinë e të dhënave, sidomos nëse ato vijnë nga burime të ndryshme. Kjo, natyrisht, është një masë e pamjaftueshme për të ruajtur integritetin e të dhënave, por në disa raste është e nevojshme. Po ashtu, të zënë kot burimet nuk është e mençur.
Pra, sensori ynë do të lancojë majat e mëtejshme të grafit nëse ka informacion të ri në postë, si dhe do të shënojë informacionin e mëparshëm si të pamjaftueshëm.
nga airflow.sensors.base_sensor_operator import BaseSensorOperator
nga airflow.utils.decorators import apply_defaults
nga my_plugin.hooks.imap_hook import IMAPHook
klasa 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 TrueMerrni dhe përdorni të dhëna
PĂ«r tĂ« marrĂ« dhe pĂ«rpunuar tĂ« dhĂ«na, mund tĂ« shkruani njĂ« operator tĂ« veçantĂ«, ose tĂ« pĂ«rdorni tĂ« gatshĂ«m. Duke qenĂ« se logjika Ă«shtĂ« akoma triviale â tĂ«rheqja e tĂ« dhĂ«nave nga emaili, pĂ«r shembull, propozoj PythonOperatorin 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)
# Konfigurimi standard i grafit
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",
)
# 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 emaili
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("Posta boshe")
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 pikave të tjera të grafit
...
# Vendos lidhjen në grafik
mail_check_sensor >> prepare_mail_data
prepare_data >> ...
# Përshkrimi i rrjedhave të tjera të menaxhimitMënyra e duhur për të organizuar postën tuaj është krijimi i një dosjeje të veçantë për emailet që dëshironi të arrini. Kjo do t'ju ndihmojë të menaxhoni informacionin në mënyrë më të lehtë. Pavarësisht sa herë që mail.ru premtuara një përmirësim në këtë drejtim, deri më sot asgjë nuk ka ndodhur. Unë kam vendosur të krijoj skedarët për mesazhet e nevojshme dhe të konfiguroj filtrat për to në ndërfaqen e postës.
Në përfundim, hapat janë si vijon: kontrolloni nëse ka emaile të reja që përmbushin kushtet, nëse po, shkarkoni arkivën nga linku në emailin e fundit.
Të dhënat e arkivës do të ruhen dhe përpunohen, dhe në fund do të kalojnë në procesin ETL, por kjo kalon jashtë tematike e artikullit. Nëse ky informacion ishte interesant dhe i dobishëm, do të jem i lumtur të vazhdoj të përshkruaj zgjidhjet ETL dhe komponentët e tyre për Apache Airflow.
Burimi: habr.com
