ETL процес за извличане на данни от електронна поща в Apache Airflow

ETL процес за извличане на данни от електронна поща в Apache Airflow

Каквото и да се случва с развитието на технологиите, винаги върви след тях дълга опашка от остарели подходи. Това може да е вследствие на плавно преминаване, човешкия фактор, технологичните нужди или нещо друго. В сферата на обработката на данни, най-показателни в това отношение са източниците на данни. Каквото и да си мечтаем да направим с това, докато част от данните се изпращат по месинджъри и имейли, да не говорим за по-архаични формати. Каня ви да разгледаме един от вариантите за Apache Airflow, който илюстрира как можем да взимаме данни от електронни писма.

Предистория

Много данни все още се предават чрез електронна поща, започвайки от междуперсонални комуникации и стигайки до стандартите за взаимодействие между компаниите. Добре е, ако можем да напишем интерфейс за получаване на данни или да сложим хора в офиса, които да въвеждат информацията в по-удобни източници, но често такава възможност просто може да липсва. Конкретната задача, пред която се изправих, е свързването на известната CRM система с хранилището на данни, а след това и с OLAP системата. Историята е такава, че за нашата компания използването на тази система беше удобно в отделна област на бизнеса. Затова всичко искане бе за възможността да оперираме с данни и от тази външна система. На първо място, разбира се, беше проучена възможността за получаване на данни от открито API. За съжаление, API-то не покриваше получаването на всички необходими данни, и, казано просто, беше в много отношения некачествено, а техническата поддръжка не пожела или не успя да отговори на искането за предоставяне на по-пълна функционалност. Затова системата предоставяше възможност за периодично получаване на липсващи данни по имейл под формата на линк за изтегляне на архив.

Важно е да се отбележи, че това не беше единственият случай, в който бизнесът искаше да събира данни от имейли или месинджъри. Въпреки това, в този случай не можехме да повлияем на външната компания, която предоставя част от данните само по този начин.

Apache Airflow

За изграждане на ETL процеси най-често използваме Apache Airflow. За да може читателят, който не е запознат с тази технология, по-добре да разбере как изглежда в контекста и изобщо, ще опиша няколко основни неща.

Apache Airflow е свободна платформа, която се използва за изграждане, изпълнение и мониторинг на ETL (Extract-Transform-Loading) процеси на езика Python. Основното понятие в Airflow е насочен ацикличен граф, където върховете на графа представляват конкретни процеси, а ребрата на графа - поток на управление или информация. Процесът може просто да извиква всяка Python функция или да има по-сложна логика с последователно извикване на няколко функции в контекста на клас. За най-честите операции вече съществуват множество готови решения, които могат да се използват като процеси. Между тях са:

  • оператори - за прехвърляне на данни от едно място на друго, например от таблица в БД в склад за данни;
  • сензори - за очакване на настъпването на определено събитие и насочване на потока на управление към следващите върхове на графа;
  • хукове - за по-нискостепенни операции, например за получаване на данни от таблица в БД (използват се в операторите);
  • и т.н.

Не е целесъобразно да описваме подробно Apache Airflow в тази статия. Кратки въведения можете да намерите тук. или тук..

Хук за получаване на данни

На първо място, за решаване на задачата трябва да напишем хук, с помощта на който да можем да:

  • съединяваме се с електронна поща;
  • намираме нужното писмо;
  • получаваме данни от писмото.

от airflow.hooks.base_hook import BaseHook
import imaplib
import logging

class IMAPHook(BaseHook):
    def __init__(self, imap_conn_id):
        """
           IMAP hook за получаване на данни от електронна поща

           :param imap_conn_id:       Идентификатор на връзката с пощата
           :type imap_conn_id:        string
        """
        self.connection = self.get_connection(imap_conn_id)
        self.mail = None

    def authenticate(self):
        """ 
            Свързваме се с пощата
        """
        mail = imaplib.IMAP4_SSL(self.connection.host)
        response, detail = mail.login(user=self.connection.login, password=self.connection.password)
        if response != "OK":
            raise AirflowException("Неуспешен вход")
        else:
            self.mail = mail

    def get_last_mail(self, check_seen=True, box="INBOX", condition="(UNSEEN)"):
        """
            Метод за получаване на идентификатора на последното писмо, 
            удовлетворяващо условията за търсене

            :param check_seen:      Маркиране на последното писмо като прочетено
            :type check_seen:       bool
            :param box:             Име на пощенската кутия
            :type box:              string
            :param condition:       Условия за търсене на писма
            :type condition:        string
        """
        self.authenticate()
        self.mail.select(mailbox=box)
        response, data = self.mail.search(None, condition)
        mail_ids = data[0].split()
        logging.info("Намерени следните писма в кутията: " + str(mail_ids))

        if not mail_ids:
            logging.info("Нови писма не са намерени")
            return None

        mail_id = mail_ids[0]

        # ако има няколко такива писма
        if len(mail_ids) > 1:
            # маркираме останалите като прочетени
            for id in mail_ids:
                self.mail.store(id, "+FLAGS", "\Seen")

            # връщаме последното
            mail_id = mail_ids[-1]

        # нужно ли е да маркираме последното като прочетено
        if not check_seen:
            self.mail.store(mail_id, "-FLAGS", "\Seen")

        return mail_id

Логиката е следната: свързваме се, намираме последното най-актуално писмо, ако има други — игнорираме ги. Използва се точно такава функция, защото по-късните писма съдържат всички данни от ранните. Ако не е така, можем да връщаме масив от всички писма или да обработим първото, а останалите — при следващото преминаване. В общи линии, всичко зависи от задачата.

Добавяме две помощни функции към хука: за изтегляне на файл и за изтегляне на файл чрез линк от имейл. Може да се помисли за извеждането им в оператор, в зависимост от честотата на използване на тази функционалност. Какво още да добавим в хука зависи от задачата: ако в имейла получаваме файлове, можем да изтеглим приложенията, ако данните пристигат в имейла, трябва да парснем имейла и т.н. В моя случай, имейлът съдържа един линк към архив, който трябва да поставя на определено място и да стартирам последващия процес на обработка.

    def download_from_url(self, url, path, chunk_size=128):
        """
            Метод за изтегляне на файл

            :param url:              Адрес на изтегляне
            :type url:               string
            :param path:             Къде да поставим файла
            :type path:              string
            :param chunk_size:       Колко байта да пишем
            :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):
        """
            Метод за изтегляне на файл по линк от имейл

            :param mail_id:         Идентификатор на имейла
            :type mail_id:          string
            :param path:            Къде да поставим файла
            :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)

Кодът е прост, така че вероятно не се нуждае от допълнителни обяснения. Ще споделя само за магическата редица imap_conn_id. Apache Airflow съхранява параметрите на връзките (логин, парола, адрес и други параметри), до които може да се получи достъп по текстов идентификатор. Визуалното управление на връзките изглежда по следния начин

ETL процес за извличане на данни от електронна поща в Apache Airflow

Сензор за очакване на данни

След като вече знаем как да се свързваме и получаваме данни от пощата, можем да напишем сензор, който да чака за тях. Не успях веднага да напиша оператор, който да обработва данните, ако те са налични, тъй като на база на получените данни от пощата работят и други процеси, включително и тези, които взимат свързани данни от други източници (API, телефония, уеб метрики и т.н.). Ще дам пример. В CRM системата се появява нов потребител, и ние все още не знаем неговото UUID. Тогава, при опит да получим данни от SIP телефонията, ще получим обаждания, свързани с неговото UUID, но не можем да ги запазим и използваме коректно. В такива случаи е важно да имаме предвид зависимостта на данните, особено ако те са от различни източници. Това, разбира се, не са достатъчни мерки за запазване на целостта на данните, но в някои случаи са необходими. А и да заемат ресурси на празно също не е рационално.

По този начин нашият сензор ще стартира последващите върхове на графа, ако има нова информация в пощата, и ще отбележи предишната информация като неактуална.

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 True

Получаваме и използваме данни

За получаване и обработка на данни може да напишете отделен оператор или да използвате готови. Тъй като логиката е тривиална — да вземем данни от имейл, за пример предлагам стандартния PythonOperator

от airflow.models импортировать DAG

от airflow.operators.python_operator импортировать PythonOperator
от airflow.sensors.my_plugin импортировать MailSensor
от my_plugin.hooks.imap_hook импортировать IMAPHook

start_date = datetime(2020, 4, 4)

# Стандартна конфигурация на графа
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",
)

# Определяме сензор
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",
)

# Функция за получаване на данни от имейл
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("Празна пощенска кутия")

    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)

# Описание на останалите възли на графа
...

# Задаваме връзка на графа
mail_check_sensor >> prepare_mail_data
prepare_data >> ...
# Описание на останалите потоци за управление

Между другото, ако вашата корпоративна поща също е на mail.ru, няма да имате достъп до търсене на писма по тема, изпращач и т.н. Те още през 2016 обещаха да го въведат, но явно са се отказали. Реших този проблем, като създадох отделна папка за нужните писма и настроих филтър в уеб интерфейса на пощата за нужните писма. По този начин в тази папка попада only необходимата поща и условията за търсене в моя случай са просто (UNSEEN).

Резюмирайки, имаме следната последователност: проверяваме дали има нови писма, съответстващи на условията, ако има, теглим архива по линка от последното писмо.
Под последните многоточия е пропуснато, че този архив ще бъде разархивиран, данните от архива ще бъдат почистени и обработени, и в края на краищата всичко това ще бъде изпратено по-нататък в конвейера на ETL процеса, но това вече излиза извън рамките на темата на статията. Ако стана интересно и полезно, с радост ще продължа да описвам ETL решения и техните части за Apache Airflow.

Източник: habr.com

Купете надежден хостинг за сайтове със защита от DDoS, VPS и VDS сървъри 🔥 Купете надежден хостинг за сайтове със защита от DDoS, VPS и VDS сървъри | ProHoster