Tworzymy pipeline do przetwarzania danych. Część 1

Cześć wszystkim. Przyjaciele, dzielimy się z wami tłumaczeniem artykułu, przygotowanym specjalnie dla studentów kursu Data Engineer. Zaczynamy!

Tworzymy pipeline do przetwarzania danych. Część 1

Apache Beam i DataFlow dla potoków w czasie rzeczywistym

Dzisiejszy post oparty jest na zadaniu, nad którym ostatnio pracowałem. Byłem naprawdę podekscytowany, mogąc je zrealizować i opisać wykonaną pracę w formacie blogowego wpisu, ponieważ dało mi to możliwość zajęcia się inżynierią danych oraz zrobienia czegoś, co byłoby bardzo przydatne dla mojego zespołu. Niedawno odkryłem, że w naszych systemach przechowywana jest dość duża ilość logów użytkowników, związanych z jednym z naszych produktów do pracy z danymi. Okazało się, że nikt nie korzystał z tych danych, więc od razu zainteresowałem się tym, co moglibyśmy się dowiedzieć, gdybyśmy zaczęli regularnie je analizować. Jednak na drodze napotkałem kilka problemów. Pierwszy problem polegał na tym, że dane były przechowywane w wielu różnych plikach tekstowych, które nie były dostępne do natychmiastowej analizy. Drugi problem polegał na tym, że były zapisane w zamkniętym systemie, więc nie mogłem użyć żadnego z moich ulubionych narzędzi do analizy danych.

Musiałem zdecydować, jak ułatwić sobie dostęp i wprowadzić jakąkolwiek wartość, integrując to źródło danych w niektóre z naszych rozwiązań do interakcji z użytkownikami. Po chwili rozmyślania postanowiłem skonstruować rurę do przesyłania tych danych do chmurowej bazy danych, aby moja drużyna mogła uzyskać do nich dostęp i zacząć generować jakiekolwiek wnioski. Po tym, jak ukończyłem specjalizację z zakresu Data Engineering na Coursera jakiś czas temu, byłem bardzo chętny, aby wykorzystać w projekcie niektóre narzędzia z kursu.

Umieszczenie danych w chmurowej bazie danych wydawało się rozsądny sposób na rozwiązanie mojego pierwszego problemu, ale co mogłem zrobić z problemem numer 2? Na szczęście istniał sposób, aby przenieść te dane do środowiska, w którym miałem dostęp do takich narzędzi, jak Python i Google Cloud Platform (GCP). Jednak był to długi proces, więc musiałem wymyślić coś, co pozwoliłoby mi kontynuować rozwój, czekając na zakończenie transferu danych. Rozwiązanie, do którego doszedłem, polegało na stworzeniu fałszywych danych przy użyciu biblioteki Faker w Pythonie. Nigdy wcześniej nie korzystałem z tej biblioteki, ale szybko zrozumiałem, jak bardzo jest przydatna. Użycie tego podejścia pozwoliło mi rozpocząć pisanie kodu i testowanie potoku bez użycia rzeczywistych danych.

Mając na uwadze to, co już powiedziałem, w tym poście opowiem, jak zbudowałem opisany wyżej potok, wykorzystując niektóre z technologii dostępnych w GCP. W szczególności będę używać Apache Beam (wersji dla Pythona), Dataflow, Pub/Sub i BigQuery do zbierania logów użytkowników, przetwarzania danych i przesyłania ich do bazy danych do dalszej analizy. W moim przypadku potrzebna mi była tylko funkcjonalność wsadowa Beam, ponieważ moje dane nie napływały w czasie rzeczywistym, więc Pub/Sub nie był wymagany. Niemniej jednak omówię wersję strumieniową, ponieważ to z tym można się spotkać w praktyce.

Wprowadzenie do GCP i Apache Beam

Google Cloud Platform oferuje zestaw naprawdę przydatnych narzędzi do przetwarzania dużych zbiorów danych. Oto niektóre z narzędzi, których będę używać:

  • Pub/Sub — to usługa komunikatów, która wykorzystuje wzorzec Wydawca-Subskrybent (Publisher-Subscriber) i pozwala nam otrzymywać dane w czasie rzeczywistym.
  • DataFlow — to serwis, który ułatwia tworzenie potoków danych i automatycznie zarządza takimi zadaniami jak skalowanie infrastruktury, co oznacza, że możemy skupić się tylko na pisaniu kodu dla naszego potoku.
  • BigQuery — to chmurowe przechowywanie danych. Jeśli znasz inne bazy danych SQL, z BigQuery nie będziesz miał problemów.
  • I w końcu będziemy używać Apache Beam, a konkretnie skupimy się na wersji Pythona do stworzenia naszego potoku. To narzędzie pozwoli nam zbudować potok do przetwarzania strumieniowego lub wsadowego, który integruje się z GCP. Jest szczególnie przydatne do przetwarzania równoległego i nadaje się do zadań typu ekstrakcji, transformacji i ładowania (ETL), więc jeśli musimy przenieść dane z jednego miejsca do drugiego, wykonując transformacje lub obliczenia, Beam jest dobrym wyborem.

Na GCP dostępnych jest wiele narzędzi, więc może być trudno śledzić je wszystkie i ich przeznaczenie, ale oto krótki podsumowanie dla odniesienia.
Na GCP dostępne jest duże ilości narzędzi, dlatego może być trudno objąć je wszystkie, w tym ich przeznaczenie, ale mimo to oto streszczenie do referencji.

Wizualizacja naszego kanału

Zobaczmy komponenty naszego kanału na rysunku 1. Na wysokim poziomie chcemy zbierać dane użytkowników w czasie rzeczywistym, przetwarzać je i przesyłać do BigQuery. Logi są tworzone, gdy użytkownicy wchodzą w interakcję z produktem, wysyłając zapytania do serwera, które następnie są rejestrowane. Te dane mogą być szczególnie przydatne do zrozumienia, jak użytkownicy współdziałają z naszym produktem i czy działają prawidłowo. Ogólnie kanał będzie zawierał następujące etapy:

Beam sprawia, że ten proces jest bardzo prosty, niezależnie od tego, czy mamy strumieniowe źródło danych, czy plik CSV, który chcemy przetworzyć w partiach. Później zobaczysz, że w kodzie występują jedynie minimalne zmiany, które są potrzebne do przełączania się między nimi. To jedna z zalet korzystania z Beama.

Tworzymy pipeline do przetwarzania danych. Część 1
Rysunek 1: Główny kanał danych: Źródło:

Tworzenie danych fikcyjnych za pomocą Faker

Jak już wspomniałem wcześniej, z powodu ograniczonego dostępu do danych postanowiłem stworzyć dane fikcyjne w tym samym formacie co rzeczywiste. To było naprawdę użyteczne ćwiczenie, ponieważ mogłem napisać kod i przetestować kanał, podczas gdy czekałem na dane. Zachęcam do zapoznania się z dokumentację Faker, jeśli chcesz dowiedzieć się, co jeszcze może zaoferować ta biblioteka. Nasze dane użytkowników będą ogólnie podobne do poniższego przykładu. Na podstawie tego formatu możemy generować dane linia po linii, aby symulować dane w czasie rzeczywistym. Te logi dostarczają nam takie informacje jak data, typ zapytania, odpowiedź serwera, adres IP itd.

192.52.197.161 - - [30/Apr/2019:21:11:42] "PUT /tag/category/tag HTTP/1.1" [401] 155 "https://harris-lopez.com/categories/about/" "Mozilla/5.0 (Macintosh; PPC Mac OS X 10_11_2) AppleWebKit/5312 (KHTML, like Gecko) Chrome/34.0.855.0 Safari/5312"

Na podstawie powyższego ciągu chcemy stworzyć naszą zmienną LINE, używając 7 zmiennych w nawiasach klamrowych poniżej. Użyjemy ich również jako nazw zmiennych w naszym schemacie tabel trochę później.

LINIA = """
{remote_addr} - - [{time_local}] "{request_type} {request_path} HTTP/1.1" [{status}] {body_bytes_sent} "{http_referer}" "{http_user_agent}"
"""

Gdybyśmy realizowali przetwarzanie wsadowe, kod byłby bardzo podobny, chociaż musielibyśmy stworzyć zestaw próbek w pewnym zakresie czasowym. Aby skorzystać z fakes, po prostu tworzymy obiekt i wywołujemy potrzebne nam metody. W szczególności Faker był przydatny do tworzenia adresów IP oraz stron internetowych. Używałem następujących metod:

fake.ipv4()
fake.uri_path()
fake.uri()
fake.user_agent()

from faker import Faker
import time
import random
import os
import numpy as np
from datetime import datetime, timedelta



LINE = """
{remote_addr} - - [{time_local}] "{request_type} {request_path} HTTP/1.1" [{status}] {body_bytes_sent} "{http_referer}" "{http_user_agent}"
"""


def generate_log_line():
    fake = Faker()
    now = datetime.now()
    remote_addr = fake.ipv4()
    time_local = now.strftime('%d/%b/%Y:%H:%M:%S')
    request_type = random.choice(["GET", "POST", "PUT"])
    request_path = "\/" + fake.uri_path()

    status = np.random.choice([200, 401, 404], p = [0.9, 0.05, 0.05])
    body_bytes_sent = random.choice(range(5, 1000, 1))
    http_referer = fake.uri()
    http_user_agent = fake.user_agent()

    log_line = LINE.format(
        remote_addr=remote_addr,
        time_local=time_local,
        request_type=request_type,
        request_path=request_path,
        status=status,
        body_bytes_sent=body_bytes_sent,
        http_referer=http_referer,
        http_user_agent=http_user_agent
    )

    return log_line

Koniec pierwszej części.

W najbliższych dniach podzielimy się z Wami kontynuacją artykułu, a teraz tradycyjnie czekamy na komentarze ;-).

Ź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