Proces ETL pozyskiwania danych z e-maila w Apache Airflow

Proces ETL pozyskiwania danych z e-maila w Apache Airflow

Bez względu na to, jak bardzo rozwijają się technologie, zawsze towarzyszy im długa lista przestarzałych podejść. Może to być spowodowane płynnością przejść, czynnikiem ludzkim, potrzebami technologicznymi lub czymś innym. W obszarze przetwarzania danych najbardziej obrazowe są źródła danych. Jak byśmy nie marzyli o pozbyciu się tego, to do momentu, w którym część danych przesyłana jest w wiadomościach i e-mailach, nie wspominając o bardziej archaicznych formatach, nie mamy innego wyboru. Zachęcam do zapoznania się z jednym z rozwiązań dla Apache Airflow, które ilustruje, jak można pozyskiwać dane z e-maili.

Tło

Wiele danych nadal przesyłanych jest za pośrednictwem e-maila, od komunikacji międzyludzkiej po standardy współpracy między firmami. Dobrze, jeśli można napisać interfejs do pozyskiwania danych lub zatrudnić ludzi w biurze, którzy wprowadzą te informacje do bardziej wygodnych źródeł, ale często taka możliwość po prostu nie istnieje. Konkretne zadanie, z którym się zmierzyłem, polegało na podłączeniu dobrze znanej systemu CRM do magazynu danych, a następnie do systemu OLAP. Tak się złożyło historycznie, że nasza firma miała wygodną pracę z tym systemem w określonym obszarze biznesowym. Dlatego wszystkim zależało na możliwości operowania danymi z tej zewnętrznej systemu. W pierwszej kolejności oczywiście zbadano możliwość pozyskiwania danych z otwartego API. Niestety, API nie obejmowało wszystkich niezbędnych danych, a mówiąc prosto, było pod wieloma względami niedopracowane, a wsparcie techniczne nie chciało ani nie mogło pomóc w udostępnieniu bardziej wyczerpującej funkcjonalności. Z drugiej strony dany system oferował możliwość okresowego otrzymywania brakujących danych na e-mail w postaci linku do pobrania archiwum.

Należy zauważyć, że to nie był jedyny przypadek, w którym biznes chciał pozyskiwać dane z wiadomości e-mail lub komunikatorów. Niemniej jednak w tym przypadku nie mogliśmy wpłynąć na zewnętrzną firmę, która dostarcza część danych tylko w taki sposób.

Apache Airflow

Do tworzenia procesów ETL najczęściej używamy Apache Airflow. Aby czytelnik nieznający tej technologii lepiej zrozumiał, jak to wygląda w kontekście i ogólnie, opiszę kilka wprowadzeń.

Apache Airflow to otwarta platforma służąca do budowy, wykonywania i monitorowania procesów ETL (Extract-Transform-Loading) w języku Python. Podstawowym pojęciem w Airflow jest skierowany acykliczny graf, w którym wierzchołki grafu to konkretne procesy, a krawędzie grafu to przepływ sterowania lub informacji. Proces może po prostu wywoływać dowolną funkcję Python lub mieć bardziej skomplikowaną logikę z sekwencyjnym wywołaniem wielu funkcji w kontekście klasy. Dla najczęściej występujących operacji dostępnych jest już wiele gotowych rozwiązań, które można wykorzystać jako procesy. Do takich rozwiązań należą:

  • operatory — do przesyłania danych z jednego miejsca do drugiego, na przykład z tabeli bazy danych do hurtowni danych;
  • sensory — do oczekiwania na wystąpienie określonego zdarzenia i kierowania przepływu sterowania do kolejnych wierzchołków grafu;
  • haki — do operacji na niższym poziomie, na przykład do pobierania danych z tabeli bazy danych (używane w operatorach);
  • itd.

Szczegółowe opisywanie Apache Airflow w tym artykule byłoby niecelowe. Krótkie wprowadzenia można znaleźć tutaj lub tutaj.

Hak do pobierania danych

Przede wszystkim, aby rozwiązać zadanie, trzeba napisać hak, za pomocą którego moglibyśmy:

  • łączyć się z pocztą elektroniczną;
  • znajdować potrzebny e-mail;
  • pobierać dane z wiadomości.

from airflow.hooks.base_hook import BaseHook
import imaplib
import logging

class IMAPHook(BaseHook):
    def __init__(self, imap_conn_id):
        """
           IMAP hook do pobierania danych z emaila

           :param imap_conn_id:       Identyfikator połączenia z pocztą
           :type imap_conn_id:        string
        """
        self.connection = self.get_connection(imap_conn_id)
        self.mail = None

    def authenticate(self):
        """ 
            Łączymy się z pocztą
        """
        mail = imaplib.IMAP4_SSL(self.connection.host)
        response, detail = mail.login(user=self.connection.login, password=self.connection.password)
        if response != "OK":
            raise AirflowException("Logowanie nie powiodło się")
        else:
            self.mail = mail

    def get_last_mail(self, check_seen=True, box="INBOX", condition="(UNSEEN)"):
        """
            Metoda do uzyskania identyfikatora ostatniego maila, 
            spełniającego warunki wyszukiwania

            :param check_seen:      Oznacz ostatni mail jako przeczytany
            :type check_seen:       bool
            :param box:             Nazwa skrzynki
            :type box:              string
            :param condition:       Warunki wyszukiwania maili
            :type condition:        string
        """
        self.authenticate()
        self.mail.select(mailbox=box)
        response, data = self.mail.search(None, condition)
        mail_ids = data[0].split()
        logging.info("W skrzynce znaleziono następujące maile: " + str(mail_ids))

        if not mail_ids:
            logging.info("Nie znaleziono nowych maili")
            return None

        mail_id = mail_ids[0]

        # jeśli jest więcej takich maili
        if len(mail_ids) > 1:
            # oznaczamy pozostałe jako przeczytane
            for id in mail_ids:
                self.mail.store(id, "+FLAGS", "\Seen")

            # zwracamy ostatni
            mail_id = mail_ids[-1]

        # czy ostatni ma być oznaczony jako przeczytany
        if not check_seen:
            self.mail.store(mail_id, "-FLAGS", "\Seen")

        return mail_id

Logika jest taka: łączymy się, znajdujemy ostatnie najnowsze e-mail, jeśli są inne — ignorujemy je. Używamy właśnie takiej funkcji, ponieważ późniejsze e-maile zawierają wszystkie dane wcześniejszych. Jeśli tak nie jest, można zwracać tablicę wszystkich e-maili lub przetwarzać pierwszy, a resztę — przy następnym przebiegu. Ogólnie rzecz biorąc, wszystko, jak zawsze, zależy od zadania.

Dodajemy do hooka dwie funkcje pomocnicze: do pobierania pliku i do pobierania pliku z linku w wiadomości e-mail. Swoją drogą, można je wyodrębnić do operatora, w zależności od częstotliwości korzystania z tej funkcjonalności. Co jeszcze dopisać do hooka, znowu zależy od zadania: jeśli w wiadomości przychodzą od razu pliki, można pobrać załączniki, jeśli dane przychodzą w wiadomości, to trzeba sparsować wiadomość itd. W moim przypadku wiadomość przychodzi z jednym linkiem do archiwum, które muszę umieścić w określonym miejscu i uruchomić dalszy proces przetwarzania.

    def download_from_url(self, url, path, chunk_size=128):
        """
            Metoda do pobierania pliku

            :param url:              Adres pobierania
            :type url:               string
            :param path:             Gdzie umieścić plik
            :type path:              string
            :param chunk_size:       Ile bajtów pisać
            :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):
        """
            Metoda do pobierania pliku z linku w wiadomości

            :param mail_id:         Identyfikator wiadomości
            :type mail_id:          string
            :param path:            Gdzie umieścić plik
            :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)

Kod jest prosty, więc raczej nie wymaga dodatkowych wyjaśnień. Powiem tylko o magicznej linii imap_conn_id. Apache Airflow przechowuje parametry połączeń (login, hasło, adres i inne parametry), do których można uzyskać dostęp za pomocą identyfikatora tekstowego. Wizualnie zarządzanie połączeniami wygląda tak

Proces ETL pozyskiwania danych z e-maila w Apache Airflow

Czujnik oczekujący na dane

Ponieważ już potrafimy łączy się i pobierać dane z poczty, teraz możemy napisać sensor do ich oczekiwania. Napisanie od razu operatora, który będzie przetwarzać dane, jeśli takie istnieją, w moim przypadku było niemożliwe, ponieważ na podstawie otrzymanych danych z poczty pracują także inne procesy, w tym biorące powiązane dane z innych źródeł (API, telefonia, metryki internetowe itp.). Podam przykład. W systemie CRM pojawił się nowy użytkownik, o którym jeszcze nie znamy UUID. Wtedy przy próbie uzyskania danych z telefonii SIP otrzymamy połączenia powiązane z jego UUID, ale nie będziemy w stanie poprawnie ich zapisać ani wykorzystać. W takich sprawach ważne jest uwzględnienie zależności danych, szczególnie jeśli pochodzą z różnych źródeł. To, oczywiście, niewystarczające środki do zachowania integralności danych, ale w niektórych przypadkach są konieczne. Zajmowanie zasobów bez sensu też jest nieefektywne.

W ten sposób nasz sensor uruchomi kolejne wierzchołki grafu, jeśli na poczcie pojawią się nowe informacje, a także oznaczy poprzednie informacje jako nieaktualne.

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

Pobieramy i wykorzystujemy dane

Aby uzyskać i przetworzyć dane, można napisać oddzielny operator lub skorzystać z gotowych. Ponieważ logika jest jak na razie trywialna — pobrać dane z wiadomości, proponuję dla przykładu standardowy PythonOperator

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)

# Standard configuration for the graph
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",
)

# Define sensor
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",
)

# Function to retrieve data from the email
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("Empty mailbox")

    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)

# Description of other graph nodes
...

# Define connection in the graph
mail_check_sensor >> prepare_mail_data
prepare_data >> ...
# Description of other control flows

Notably, if your corporate email is also on mail.ru, you won't have access to search emails by subject, sender, etc. They promised to implement this back in 2016, but apparently changed their minds. I solved this issue by creating a separate folder for the needed emails and setting up a filter for the required emails in the webmail interface. Thus, only the needed emails go into this folder, and the search condition in my case is simply (UNSEEN).

In summary, we have the following sequence: check if there are new emails that meet the criteria, and if so, download the archive from the link in the last email.
Under the last ellipses, it is omitted that this archive will be unpacked, the data from the archive will be cleared and processed, and ultimately this will continue on to the ETL process pipeline, but this is beyond the scope of this article. If you found this interesting and useful, I would be happy to continue describing ETL solutions and their components for Apache Airflow.

Źródło: habr.com

Kup solidny hosting stron z ochroną przed DDoS, serwery VPS VDS 🔥 Kup solidny hosting stron z ochroną przed DDoS, serwery VPS VDS | ProHoster