Tworzymy potok przetwarzania danych. Część 2

Cześć wszystkim. Dzielimy się tłumaczeniem ostatniej części artykułu, przygotowanego specjalnie dla studentów kursu Data Engineer. Z pierwszą częścią można zapoznać się tutaj.

Apache Beam i DataFlow dla potoków w czasie rzeczywistym

Tworzymy potok przetwarzania danych. Część 2

Konfiguracja Google Cloud

Uwaga: Aby uruchomić potok i publikować dane logu użytkownika, użyłem Google Cloud Shell, ponieważ miałem problemy z uruchomieniem potoku w Pythonie 3. Google Cloud Shell korzysta z Pythona 2, co lepiej współpracuje z Apache Beam.

Aby uruchomić potok, musimy trochę pogrzebać w ustawieniach. Ci z was, którzy wcześniej nie korzystali z GCP, muszą wykonać następujące 6 kroków, które są przedstawione na tej stronie.

Następnie musimy załadować nasze skrypty do chmury Google i skopiować je do naszej Google Cloud Shell. Ładowanie do chmury jest dość trywialne (opis można znaleźć tutaj). Aby skopiować nasze pliki, możemy otworzyć Google Cloud Shell z paska narzędzi, klikając pierwszy ikonę po lewej stronie na rysunku 2 poniżej.

Tworzymy potok przetwarzania danych. Część 2
Rysunek 2

Polecenia, których potrzebujemy do skopiowania plików i zainstalowania niezbędnych bibliotek, są wymienione poniżej.

# Copy file from cloud storage
gsutil cp gs://<YOUR-BUCKET>/ * .
sudo pip install apache-beam[gcp] oauth2client==3.0.0
sudo pip install -U pip
sudo pip install Faker==1.0.2
# Environment variables
BUCKET=<YOUR-BUCKET>
PROJECT=<YOUR-PROJECT>

Tworzenie naszej bazy danych i tabeli

Po wykonaniu wszystkich kroków związanych z konfiguracją, następną rzeczą, którą musimy zrobić, jest stworzenie zestawu danych i tabeli w BigQuery. Istnieje kilka sposobów na zrobienie tego, ale najłatwiejszym jest użycie konsoli Google Cloud, zaczynając od utworzenia zestawu danych. Możesz wykonać kroki opisane poniżej linkiem, aby stworzyć tabelę ze schematem. Nasza tabela będzie miała 7 kolumn, które odpowiadają komponentom każdego logu użytkownika. Dla wygody zdefiniujemy wszystkie kolumny jako ciągi (typ string), z wyjątkiem zmiennej timelocal, i nazwiemy je zgodnie z wcześniej wygenerowanymi zmiennymi. Schemat naszej tabeli powinien wyglądać jak na rysunku 3.

Tworzymy potok przetwarzania danych. Część 2
Rysunek 3. Schemat tabeli

Publikacja danych logu użytkownika

Pub/Sub jest krytycznym komponentem naszego pipeline'u, ponieważ pozwala na interakcję między niezależnymi aplikacjami. Działa jako pośrednik, umożliwiając nam wysyłanie i odbieranie wiadomości pomiędzy aplikacjami. Pierwszym krokiem, który musimy wykonać, jest utworzenie tematu (topic). Wystarczy przejść do Pub/Sub w konsoli i nacisnąć CREATE TOPIC.

Poniższy kod wywołuje nasz skrypt generujący dane logów, określone powyżej, a następnie łączy się i wysyła logi do Pub/Sub. Jedynym, co musimy zrobić, jest stworzenie obiektu PublisherClient, wskazanie ścieżki do tematu za pomocą metody topic_path i wywołanie funkcji publish z topic_path z danymi. Zauważ, że importujemy generate_log_line z naszego skryptu stream_logs, więc upewnij się, że te pliki znajdują się w tym samym folderze, w przeciwnym razie otrzymasz błąd importu. Następnie możemy uruchomić to przez naszą konsolę Google, używając:

python publish.py

from stream_logs import generate_log_line
import logging
from google.cloud import pubsub_v1
import random
import time


PROJECT_ID="user-logs-237110"
TOPIC = "userlogs"


publisher = pubsub_v1.PublisherClient()
topic_path = publisher.topic_path(PROJECT_ID, TOPIC)

def publish(publisher, topic, message):
    data = message.encode('utf-8')
    return publisher.publish(topic_path, data = data)

def callback(message_future):
    # Kiedy timeout jest niesprecyzowany, metoda wyjątku czeka bez końca.
    if message_future.exception(timeout=30):
        print('Publikacja wiadomości na {} spowodowała wyjątek {}.'.format(
            topic_name, message_future.exception()))
    else:
        print(message_future.result())


if __name__ == '__main__':

    while True:
        line = generate_log_line()
        print(line)
        message_future = publish(publisher, topic_path, line)
        message_future.add_done_callback(callback)

        sleep_time = random.choice(range(1, 3, 1))
        time.sleep(sleep_time)

Gdy plik zostanie uruchomiony, będziemy mogli zobaczyć dane logu w konsoli, jak pokazano na poniższym obrazku. Ten skrypt będzie działał, dopóki nie użyjemy CTRL+C, aby go zakończyć.

Tworzymy potok przetwarzania danych. Część 2
Rysunek 4. Wyjście publish_logs.py

Pisanie kodu naszego pipeline'u

Teraz, gdy wszystko jest gotowe, możemy przejść do najciekawszej części — napisania kodu naszego pipeline'u z użyciem Beam i Pythona. Aby utworzyć pipeline Beam, musimy stworzyć obiekt pipeline (p). Po utworzeniu obiektu pipeline możemy zastosować kilka funkcji jedna po drugiej, używając operatora pipe (|). Ogólnie, przepływ pracy wygląda jak na poniższym obrazku.

[Final Output PCollection] = ([Initial Input PCollection] | [First Transform]
             | [Second Transform]
             | [Third Transform])

W naszym kodzie stworzymy dwie funkcje użytkownika. Funkcję regex_clean, która skanuje dane i wyciąga odpowiedni ciąg na podstawie listy PATTERNS, wykorzystując funkcję re.search. Funkcja zwraca ciąg rozdzielony przecinkami. Jeśli nie jesteś ekspertem w zakresie wyrażeń regularnych, zalecam zapoznanie się z tym samouczkiem i poćwiczenie w notatniku, aby przetestować kod. Następnie definiujemy użytkownikowską funkcję ParDo o nazwie Split, która jest wariacją przekształcenia Beam do przetwarzania równoległego. W Pythonie robi się to w szczególny sposób — musimy stworzyć klasę, która dziedziczy z klasy DoFn Beam. Funkcja Split przyjmuje sparsowany ciąg z poprzedniej funkcji i zwraca listę słowników z kluczami odpowiadającymi nazwom kolumn w naszej tabeli BigQuery. Jest jedna rzecz, którą warto zauważyć na temat tej funkcji: musiałem zaimportować datetime wewnątrz funkcji, aby działała. Otrzymywałem komunikat o błędzie przy imporcie na początku pliku, co było dziwne. Ta lista jest następnie przekazywana do funkcji WriteToBigQuery, która po prostu dodaje nasze dane do tabeli. Kod dla Batch DataFlow Job i Streaming DataFlow Job przedstawiono poniżej. Jedyna różnica między kodem pakietowym a strumieniowym polega na tym, że w przetwarzaniu pakietów odczytujemy CSV z src_path, używając funkcji ReadFromText z Beam.

Batch DataFlow Job (przetwarzanie pakietów)

import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions
from google.cloud import bigquery
import re
import logging
import sys

PROJECT='user-logs-237110'
schema = 'remote_addr:STRING, timelocal:STRING, request_type:STRING, status:STRING, body_bytes_sent:STRING, http_referer:STRING, http_user_agent:STRING'


src_path = "user_log_fileC.txt"

def regex_clean(data):

    PATTERNS =  [r'(^S+.[S+.]+S+)s',r'(?<=[).+?(?=])',
           r'"(S+)s(S+)s*(S*)"',r's(d+)s',r"(?> beam.io.textio.ReadFromText(src_path)
      | "clean address" >> beam.Map(regex_clean)
      | 'ParseCSV' >> beam.ParDo(Split())
      | 'WriteToBigQuery' >> beam.io.WriteToBigQuery('{0}:userlogs.logdata'.format(PROJECT), schema=schema,
        write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND)
   )

   p.run()

if __name__ == '__main__':
  logger = logging.getLogger().setLevel(logging.INFO)
  main()

Streaming DataFlow Job (przetwarzanie strumieniowe)

from apache_beam.options.pipeline_options import PipelineOptions
from google.cloud import pubsub_v1
from google.cloud import bigquery
import apache_beam as beam
import logging
import argparse
import sys
import re


PROJECT="user-logs-237110"
schema = 'remote_addr:STRING, timelocal:STRING, request_type:STRING, status:STRING, body_bytes_sent:STRING, http_referer:STRING, http_user_agent:STRING'
TOPIC = "projects/user-logs-237110/topics/userlogs"


def regex_clean(data):

    PATTERNS =  [r'(^S+.[S+.]+S+)s',r'(?<=[).+?(?=])',
           r'"(S+)s(S+)s*(S*)"',r's(d+)s',r"(?> beam.io.ReadFromPubSub(topic=TOPIC).with_output_types(bytes)
      | "Decode" >> beam.Map(lambda x: x.decode('utf-8'))
      | "Clean Data" >> beam.Map(regex_clean)
      | 'ParseCSV' >> beam.ParDo(Split())
      | 'WriteToBigQuery' >> beam.io.WriteToBigQuery('{0}:userlogs.logdata'.format(PROJECT), schema=schema,
        write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND)
   )
   result = p.run()
   result.wait_until_finish()

if __name__ == '__main__':
  logger = logging.getLogger().setLevel(logging.INFO)
  main()

Uruchamianie potoku

Możemy uruchomić potok na kilka różnych sposobów. Gdybyśmy chcieli, moglibyśmy po prostu uruchomić go lokalnie z terminala, logując się zdalnie do GCP.

python -m main_pipeline_stream.py 
 --input_topic "projects/user-logs-237110/topics/userlogs" 
 --streaming

Jednak zamierzamy uruchomić go za pomocą DataFlow. Możemy to zrobić za pomocą poniższej komendy, podając następujące obowiązkowe parametry.

  • projekt — ID swojego projektu GCP.
  • runner — narzędzie uruchamiania potoku, które przeanalizuje Twój program i skonstruuje Twój potok. Aby uruchomić w chmurze, musisz wskazać DataflowRunner.
  • staging_location — ścieżka do Cloud Storage Dataflow, w celu indeksowania pakietów kodu wymaganych przez przetwarzające zadanie.
  • temp_location — ścieżka do Cloud Storage Dataflow do przechowywania tymczasowych plików zadań tworzonych podczas uruchamiania potoku.
  • streaming

python main_pipeline_stream.py 
--runner DataFlow 
--project $PROJECT 
--temp_location $BUCKET/tmp 
--staging_location $BUCKET/staging
--streaming

Gdy to polecenie jest wykonywane, możemy przejść do zakładki DataFlow w konsoli Google i zobaczyć nasz pipeline. Klikając na pipeline, powinniśmy zobaczyć coś podobnego do rysunku 4. Dla celów debugowania może być bardzo pomocne przejście do logów, a następnie do Stackdriver, aby zobaczyć szczegółowe logi. Pomogło mi to rozwiązać problemy z pipeline'em w wielu przypadkach.

Tworzymy potok przetwarzania danych. Część 2
Rysunek 4: Pipeline Beam

Dostęp do naszych danych w BigQuery

Zatem nasz pipeline powinien już działać i przesyłać dane do naszej tabeli. Aby to sprawdzić, możemy przejść do BigQuery i zobaczyć dane. Po użyciu poniższego polecenia powinieneś zobaczyć pierwsze kilka wierszy zestawu danych. Teraz, gdy mamy dane przechowywane w BigQuery, możemy przeprowadzić dalszą analizę, a także podzielić się danymi z kolegami i zacząć odpowiadać na pytania biznesowe.

SELECT * FROM `user-logs-237110.userlogs.logdata` LIMIT 10;

Tworzymy potok przetwarzania danych. Część 2
Rysunek 5: BigQuery

Podsumowanie

Mamy nadzieję, że ten post będzie przydatnym przykładem tworzenia strumieniowego pipeline'u danych oraz poszukiwania sposobów na uczynienie danych bardziej dostępnymi. Przechowywanie danych w takim formacie daje nam wiele korzyści. Teraz możemy zacząć odpowiadać na ważne pytania, takie jak, ile osób korzysta z naszego produktu? Czy baza użytkowników rośnie w czasie? Z jakimi aspektami produktu użytkownicy najczęściej wchodzą w interakcję? I czy są błędy tam, gdzie nie powinno ich być? To pytania, które będą interesujące dla organizacji. Na podstawie spostrzeżeń, które wynikają z odpowiedzi na te pytania, możemy udoskonalić produkt i zwiększyć zaangażowanie użytkowników.

Beam jest rzeczywiście przydatny do tego typu ćwiczeń, a także ma wiele innych interesujących zastosowań. Na przykład, możesz analizować dane dotyczące ticków giełdowych w czasie rzeczywistym i podejmować decyzje handlowe na podstawie analizy. Możliwe, że masz dane z czujników z pojazdów i chcesz obliczyć poziom natężenia ruchu. Możesz również być firmą gier, która zbiera dane o użytkownikach i wykorzystuje je do tworzenia pulpitów nawigacyjnych do śledzenia kluczowych wskaźników. Dobrze, panowie, to już temat na inny post, dziękuję za przeczytanie, a dla tych, którzy chcą zobaczyć pełny kod, poniżej znajduje się link do mojego GitHuba.

https://github.com/DFoly/User_log_pipeline

Na tym kończymy. Przeczytaj pierwszą część.

Ź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