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 импортировать BaseHook
импортировать imaplib
импортировать logging

класс IMAPHook(BaseHook):
    def __init__(self, imap_conn_id):
        """
           IMAP хук для получения данных с электронной почты

           :param imap_conn_id:       Идентификатор подключения к почте
           :type imap_conn_id:        строка
        """
        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)
        если response != "OK":
            raise AirflowException("Не удалось войти")
        еще:
            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:              строка
            :param condition:       Условия поиска писем
            :type condition:        строка
        """
        self.authenticate()
        self.mail.select(mailbox=box)
        response, data = self.mail.search(None, condition)
        mail_ids = data[0].split()
        logging.info("В ящике найдены следующие письма: " + str(mail_ids))

        если не mail_ids:
            logging.info("Новые письма не найдены")
            вернуть None

        mail_id = mail_ids[0]

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

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

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

        вернуть 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, но явно се отказаха. Реших този проблем, като създадох отделна папка за нужните имейли и настроих филтър в уеб интерфейса на пощата. Така само нужните имейли попадат в тази папка и условията за търсене при мен са прости (UNSEEN).

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

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

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