Creăm un pipeline pentru procesarea datelor în flux. Parte 2

Bună tuturor. Împărtășim traducerea părții finale a articolului, pregătit special pentru studenții cursului „Inginer de date”. Puteți consulta prima parte aici.

Apache Beam și DataFlow pentru conducte în timp real

Creăm un pipeline pentru procesarea datelor în flux. Parte 2

Configurarea Google Cloud

Notă: Pentru a rula pipeline-ul și a publica datele jurnalului personalizat, am folosit Google Cloud Shell, deoarece am întâmpinat probleme cu rularea pipeline-ului pe Python 3. Google Cloud Shell folosește Python 2, care este mai compatibil cu Apache Beam.

Pentru a rula pipeline-ul, trebuie să ne uităm puțin la setări. Pentru cei dintre voi care nu ați folosit GCP înainte, este necesar să urmați următorii 6 pași prezentati pe această apărea un mesaj că laserul nu este disponibil pentru moment și un indiciu: în birou s-au fumat siguranțele, trebuie să suni compania de administrare și să ceri alimentarea..

După aceea, va trebui să încărcăm scripturile noastre în Google Cloud Storage și să le copiem în Google Cloud Shell. Încărcarea în Google Cloud Storage este destul de simplă (informații pot fi găsite aici). Pentru a copia fișierele noastre, putem deschide Google Cloud Shell din bara de instrumente, făcând clic pe primul simbol din stânga în figura 2 de mai jos.

Creăm un pipeline pentru procesarea datelor în flux. Parte 2
Figura 2

Comenzile de care avem nevoie pentru a copia fișierele și a instala bibliotecile necesare sunt enumerate mai jos.

# 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>

Crearea bazei noastre de date și a tabelului

După ce am parcurs toți pașii necesari pentru configurare, următorul lucru pe care trebuie să-l facem este să creăm un set de date și un tabel în BigQuery. Există mai multe moduri de a face acest lucru, dar cel mai simplu este să folosim consola Google Cloud, creând mai întâi un set de date. Puteți urma pașii indicați mai jos linkul, pentru a crea un tabel cu un schema. Tabelul nostru va avea 7 coloane, corespunzătoare componentelor fiecărui jurnal personalizat. Pentru comoditate, vom defini toate coloanele ca șiruri (tip string), cu excepția variabilei timelocal, și le vom numi conform variabilelor pe care le-am generat anterior. Schema tabelului nostru ar trebui să arate ca în figura 3.

Creăm un pipeline pentru procesarea datelor în flux. Parte 2
Figura 3. Schema tabelului

Publicarea datelor jurnalului personalizat

Pub/Sub este un component critic al pipeline-ului nostru, deoarece permite mai multor aplicații independente să interacționeze între ele. În special, funcționează ca un intermediar, permițându-ne să trimitem și să primim mesaje între aplicații. Primul lucru pe care trebuie să-l facem este să creăm un subiect (topic). Este suficient să mergem la Pub/Sub în consolă și să facem clic pe CREATE TOPIC.

Codul de mai jos apelează scriptul nostru pentru generarea datelor de jurnal definite mai sus și apoi se conectează și trimite jurnalele în Pub/Sub. Tot ce trebuie să facem este să creăm un obiect PublisherClient, să specificăm calea către subiect folosind metoda topic_path și să apelăm funcția publish de topic_path și datele. Observați că importăm generate_log_line din scriptul nostru stream_logs, așa că asigurați-vă că aceste fișiere se află în aceeași folder, altfel veți primi o eroare de import. Apoi, putem rula aceasta prin consola noastră Google, utilizând:

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):
    # Când timeout-ul nu este specificat, metoda de excepție așteaptă indefinite.
    if message_future.exception(timeout=30):
        print('Publicarea mesajului în {} a generat o Excepție {}.'.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)

Odată ce fișierul va fi pornit, vom putea observa ieșirea datelor de jurnal pe consolă, așa cum este arătat în imaginea de mai jos. Acest script va funcționa până când vom utiliza CTRL+C, pentru a-l finaliza.

Creăm un pipeline pentru procesarea datelor în flux. Parte 2
Figura 4. Ieșirea publish_logs.py

Scrierea codului pentru canalul nostru

Acum, că am pregătit totul, putem trece la partea cea mai interesantă — scrierea codului pentru canalul nostru, folosind Beam și Python. Pentru a crea un canal Beam, trebuie să creăm un obiect canal (p). După ce am creat obiectul canal, putem aplica mai multe funcții unele după altele, folosind operatorul pipe (|). În general, fluxul de lucru arată ca în imaginea de mai jos.

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

În codul nostru, vom crea două funcții personalizate. Funcția regex_clean, care scanează datele și extrage șirul corespunzător pe baza listei PATTERNS, folosind funcția re.search. Funcția returnează un șir despărțit prin virgulă. Dacă nu ești expert în expresii regulate, îți recomand să te familiarizezi cu acest tutorial și să exersez în notepad pentru a verifica codul. După aceasta, definim o funcție ParDo personalizată numită Split, care este o variație a transformării Beam pentru procesarea paralelă. În Python, acest lucru se face într-un mod special — trebuie să creăm o clasă care moștenește clasa DoFn Beam. Funcția Split ia un șir analizat din funcția anterioară și returnează o listă de dicționare cu chei corespunzătoare numelui coloanelor din tabela noastră BigQuery. Este ceva de notat despre această funcție: a trebuit să import funcția datetime în interiorul funcției pentru ca aceasta să funcționeze. Am primit un mesaj de eroare la importul din începutul fișierului, ceea ce era ciudat. Această listă este apoi transmisă funcției WriteToBigQuery, care pur și simplu adaugă datele noastre în tabel. Codul pentru Batch DataFlow Job și Streaming DataFlow Job este prezentat mai jos. Singura diferență între codul de procesare în lot și cel în flux este că în procesarea în lot citim CSV din src_path, folosind funcția ReadFromText din Beam.

Batch DataFlow Job (procesare în lot)

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 (procesare în flux)

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()

Pornirea conductorului

Putem să pornim conductorul în mai multe moduri. Dacă am dori, am putea să-l rulăm local din terminal, conectându-ne de la distanță la GCP.

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

Totuși, vom rula folosind DataFlow. Putem face acest lucru cu comanda de mai jos, setând următoarele opțiuni obligatorii.

  • project — ID-ul proiectului vostru GCP.
  • runner — instrumentul de lansare a conductorului, care va analiza programul vostru și va construi conductorul. Pentru a rula în cloud, trebuie să specificați DataflowRunner.
  • staging_location — calea către stocarea Cloud Dataflow pentru indexarea pachetelor de cod necesare procesatorilor care efectuează sarcina.
  • temp_location — calea către stocarea Cloud Dataflow pentru stocarea fișierelor temporare de sarcini create în timpul execuției conductorului.
  • streaming

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

Cât timp această comandă este executată, putem să mergem la tab-ul DataFlow în consola Google pentru a vizualiza conducta noastră. Facând clic pe conducta, ar trebui să vedem ceva asemănător cu figura 4. În scopuri de depanare, poate fi foarte util să accesăm jurnalele, apoi Stackdriver pentru a vizualiza jurnalele detaliate. Acest lucru m-a ajutat să rezolv probleme cu conducta în mai multe cazuri.

Creăm un pipeline pentru procesarea datelor în flux. Parte 2
Figura 4: Conducta Beam

Acces la datele noastre în BigQuery

Deci, ar trebui să avem deja o conductă activă cu datele care intră în tabelul nostru. Pentru a verifica acest lucru, putem să mergem în BigQuery și să vizualizăm datele. După utilizarea comenzii de mai jos, ar trebui să vedeți primele câteva rânduri ale setului de date. Acum, când avem datele stocate în BigQuery, putem efectua analize ulterioare și putem împărtăși datele cu colegii, precum și să începem să răspundem la întrebările de afaceri.

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

Creăm un pipeline pentru procesarea datelor în flux. Parte 2
Figura 5: BigQuery

Concluzie

Sperăm că această postare va fi un exemplu util de creare a unei conducte de date în flux, precum și de găsire a modalităților de a face datele mai accesibile. Stocarea datelor într-un astfel de format ne oferă multe avantaje. Acum putem începe să răspundem la întrebările importante, cum ar fi câți oameni folosesc produsul nostru? Crește baza de utilizatori în timp? Cu ce aspecte ale produsului interacționează cel mai mult utilizatorii? Și există erori acolo unde nu ar trebui? Acestea sunt întrebările care vor interesa organizația. Pe baza ideilor ce decurg din răspunsurile la aceste întrebări, vom putea îmbunătăți produsul și crește implicarea utilizatorilor.

Beam este cu adevărat util pentru acest tip de exerciții, având totodată o serie de alte cazuri de utilizare interesante. De exemplu, poți analiza datele de pe piața bursieră în timp real și să faci tranzacții pe baza analizei; poate ai date de la senzori care vin de la vehicule și vrei să calculezi nivelul de trafic. De asemenea, poți fi o companie de jocuri care colectează date despre utilizatori și le folosește pentru a crea panouri de control pentru a urmări indicatorii cheie. Ei bine, domnilor, aceasta este o temă pentru un alt post, mulțumesc pentru lectură, iar pentru cei care doresc să vadă codul complet, mai jos este linkul către GitHub-ul meu.

https://github.com/DFoly/User_log_pipeline

Asta e tot. Citește prima parte.

Sursa: habr.com

Cumpără un hosting fiabil pentru site-uri cu protecție DDoS, servere VPS VDS 🔥 Cumpără un hosting fiabil pentru site-uri cu protecție DDoS, servere VPS VDS | ProHoster