Cześć wszystkim. Dzielimy się tłumaczeniem ostatniej części artykułu, przygotowanego specjalnie dla studentów kursu . Z pierwszą częścią można zapoznać się .
Apache Beam i DataFlow dla potoków w czasie rzeczywistym

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 .
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źć ). 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.

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 , 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.

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.pyfrom 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ć.

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 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"
--streamingJednak 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.

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; 
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.
Na tym kończymy. .
Źródło: habr.com
