Krijojmë një kanal për përpunimin e të dhënave. Pjesa 2

Përshëndetje të gjithëve. Po ndajmë përkthimin e pjesës përfundimtare të artikullit, përgatitur posaçërisht për studentët e kursit Inxhinier i të Dhënave. Mund të njiheni me pjesën e parë këtu.

Apache Beam dhe DataFlow për pipeline të vërtetë në kohë reale

Krijojmë një kanal për përpunimin e të dhënave. Pjesa 2

Konfigurimi i Google Cloud

Shënim: Për të ekzekutuar linjën e prodhimit dhe publikimin e të dhënave nga regjistri i përdoruesit, kam përdorur Google Cloud Shell, pasi kam pasur probleme me ekzekutimin e linjës në Python 3. Google Cloud Shell përdor Python 2, i cili është më i përshtatshëm për Apache Beam.

Për të ekzekutuar linjën, na nevojitet pak të thellojmë në konfigurime. Për ata prej jush që nuk e kanë përdorur më parë GCP, duhet të kryeni 6 hapat e mëposhtëm që janë cituar në këtë faqja.

Pas kësaj, do të duhet të ngarkojmë skenarët tanë në ruajtjen e re të Google dhe t'i kopjojmë në Google Cloud Shell tonë. Ngarkimi në ruajtjen e re është mjaft i thjeshtë (përshkrimi mund të gjendet këtu). Për të kopjuar skedarët tanë, mund të hapim Google Cloud Shell nga paneli i mjeteve duke klikuar ikonen e parë nga majtas në figurën 2 më poshtë.

Krijojmë një kanal për përpunimin e të dhënave. Pjesa 2
Figura 2

Komandat që na nevojiten për kopjimin e skedarëve dhe instalimin e bibliotekave të nevojshme janë listuar më poshtë.

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

Krijimi i bazës sonë të të dhënave dhe tabelës

Pas kryerjes së të gjitha hapave të lidhura me konfigurimin, gjëja tjetër që na nevojitet të bëjmë është krijimi i një grupi të dhënash dhe tabelës në BigQuery. Ka disa mënyra për ta bërë këtë, por më e thjeshta është përdorimi i konsolës së Google Cloud, duke krijuar fillimisht një grup të dhënash. Ju mund të kryeni veprimet e listuara nga linkun, për të krijuar një tabelë me skemë. Tabela jonë do të ketë 7 kolona, të cilat korrespondojnë me komponentët e çdo regjistri të përdoruesit. Për shkak të lehtësisë, ne do t'i përcaktojmë të gjitha kolonat si string (tipi string), përveç variablës timelocal, dhe do t'i quajmë sipas variablave që kemi gjeneruar më parë. Skema e tabelës sonë duhet të duket si në figurën 3.

Krijojmë një kanal për përpunimin e të dhënave. Pjesa 2
Figura 3. Skema e tabelës

Publikimi i të dhënave nga regjistri i përdoruesit

Pub/Sub është një komponent kritik i linjës tonë, pasi lejon aplikacione të pavarura të komunikojnë midis tyre. Në veçanti, ai funksionon si një ndërmjetës, duke na lejuar të dërgojmë dhe marrim mesazhe midis aplikacioneve. E para që na nevojitet të bëjmë është krijimi i një teme (topic). Mjafton të kaloni në Pub/Sub në konsolë dhe të klikoni krijo teme (CREATE TOPIC).

Kodi i mëposhtëm thërret skriptin tonë për të gjeneruar të dhënat e logutsi, të cilat janë përcaktuar më parë, dhe pastaj lidhet dhe dërgon loget në Pub/Sub. E vetmja gjë që na duhet të bëjmë është të krijojmë një objekt PublisherClient, të përcaktojmë rrugën e temës me metodën topic_path dhe të thërrasim funksionin publish me topic_path me të dhënat. Vini re se ne importojmë generate_log_line nga skripti ynë stream_logs, prandaj sigurohuni që këto skedarë janë në të njëjtën dosje, përndryshe do të merrni një gabim importi. Më pas, mund ta ekzekutojmë këtë përmes konsolës sonë të Google, duke përdorur:

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):
    # Kur koha e skadencës nuk është e përcaktuar, metoda e përjashtimit presin pafundësisht.
    if message_future.exception(timeout=30):
        print('Dërgimi i mesazhit në {} ndodhi një përjashtim {}.'.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)

Sa kohë që skedari është i aktivizuar, do të mund të shohim daljen e të dhënave të logut në konsolë, siç tregohet në figurën më poshtë. Ky skript do të funksionojë derisa të përdorim CTRL+C, për ta ndaluar atë.

Krijojmë një kanal për përpunimin e të dhënave. Pjesa 2
Figura 4. Dalja publish_logs.py

Shkrimi i kodit tonë për tubin

Tani qĂ« kemi pĂ«rgatitur gjithçka, mund tĂ« kalojmĂ« nĂ« pjesĂ«n mĂ« interesante — shkrimi i kodit tonĂ« pĂ«r tubin, duke pĂ«rdorur Beam dhe Python. PĂ«r tĂ« krijuar njĂ« tub tĂ« Beam-it, na nevojitet tĂ« krijojmĂ« njĂ« objekt tubi (p). Pasi tĂ« kemi krijuar objektin e tubit, mund tĂ« aplikojmĂ« disa funksione njĂ«ra pas tjetrĂ«s, duke pĂ«rdorur operatorin pipe (|). NĂ« pĂ«rgjithĂ«si, procesi i punĂ«s duket si nĂ« figurĂ«n mĂ« poshtĂ«.

[Dalja përfundimtare PCollection] = ([PCollection e Input të Parë] | [Transformimi i Parë]
             | [Transformimi i Dytë]
             | [Transformimi i Tretë])

NĂ« kodin tonĂ« do tĂ« krijojmĂ« dy funksione tĂ« personalizuara. Funksionin regex_clean, i cili skanon tĂ« dhĂ«nat dhe nxjerr njĂ« varg tĂ« pĂ«rshtatshĂ«m nĂ« bazĂ« tĂ« listĂ«s PATTERNS, duke pĂ«rdorur funksionin re.search. Funksioni kthen njĂ« varg tĂ« ndarĂ« me virgula. NĂ«se nuk jeni ekspert nĂ« shprehjet e rregullta, ju rekomandoj tĂ« njiheni me kĂ«tĂ« tutorial dhe tĂ« praktikoni nĂ« blloknot pĂ«r tĂ« verifikuar kodin. Pas kĂ«saj, ne pĂ«rcaktojmĂ« njĂ« funksion ParDo tĂ« personalizuar tĂ« quajtur Ndarja, i cili Ă«shtĂ« njĂ« variacion i transformimit Beam pĂ«r pĂ«rpunim paralel. NĂ« Python, kjo bĂ«het nĂ« njĂ« mĂ«nyrĂ« tĂ« veçantĂ« — ne duhet tĂ« krijojmĂ« njĂ« klasĂ« qĂ« trashĂ«gon nga klasa DoFn e Beam. Funksioni Ndarja merr njĂ« varg tĂ« analizuar nga funksioni paraprak dhe kthen njĂ« listĂ« dictionaries me çelĂ«sa qĂ« korrespondojnĂ« me emrat e kolonave nĂ« tabelĂ«n tonĂ« BigQuery. Ka diçka qĂ« duhet theksuar nĂ« lidhje me kĂ«tĂ« funksion: mĂ« duhej tĂ« importoja datetime brenda funksionit pĂ«r tĂ« punuar. Marrja e njĂ« mesazhi gabimi gjatĂ« importimit nĂ« fillim tĂ« skedarit ishte e çuditshme. Kjo listĂ« mĂ« pas i kalon funksionit WriteToBigQuery, i cili thjesht shton tĂ« dhĂ«nat tona nĂ« tabelĂ«. Kodi pĂ«r Batch DataFlow Job dhe Streaming DataFlow Job Ă«shtĂ« i dhĂ«nĂ« mĂ« poshtĂ«. Diferenca e vetme midis kodit paketor dhe atij rrjedhĂ«sor Ă«shtĂ« se nĂ« pĂ«rpunimin e paketave lexojmĂ« CSV nga src_path, duke pĂ«rdorur funksionin ReadFromText nga Beam.

Batch DataFlow Job (përpunimi i paketave)

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"(?<=[).d+(?=])",
           r'"[A-Z][a-z]+', r'"(http|https):\/\/[a-z]+.[a-z]+.[a-z]+']
    result = []
    for match in PATTERNS:
      try:
        reg_match = re.search(match, data).group()
        if reg_match:
          result.append(reg_match)
        else:
          result.append(" ")
      except:
        print("Ka ndodhur një gabim me kërkimin regex")
    result = [x.strip() for x in result]
    result = [x.replace('"', "") for x in result]
    res = ','.join(result)
    return res


class Split(beam.DoFn):

    def process(self, element):
        from datetime import datetime
        element = element.split(",")
        d = datetime.strptime(element[1], "%d\/\b\/\%Y:%H:%M:%S")
        date_string = d.strftime("%Y-%m-%d %H:%M:%S")

        return [{ 
            'remote_addr': element[0],
            'timelocal': date_string,
            'request_type': element[2],
            'status': element[3],
            'body_bytes_sent': element[4],
            'http_referer': element[5],
            'http_user_agent': element[6]
    
        }]

def main():

   p = beam.Pipeline(options=PipelineOptions())

   (p
      | 'ReadData' >> beam.io.textio.ReadFromText(src_path)
      | "pastroni adresën" >> 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 (përpunimi i rrjedhës)

nga apache_beam.options.pipeline_options import PipelineOptions
nga google.cloud import pubsub_v1
nga google.cloud import bigquery
import apache_beam si beam
import logging
import argparse
import sys
import re


PROJEKTI="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'
TOPIKU = "projekti/user-logs-237110/topikët/userlogs"


def regex_clean(të dhëna):

    PATTERNET =  [r'(^S+.[S+.]+S+)s',r'(?<=[).+?(?=])',
           r'"(S+)s(S+)s*(S*)"',r's(d+)s',r"(?> beam.io.ReadFromPubSub(topic=TOPIKU).with_output_types(bytes)
      | "Decode" >> beam.Map(lambda x: x.decode('utf-8'))
      | "Pastrimi i të Dhënave" >> beam.Map(regex_clean)
      | 'AnalizaCSV' >> beam.ParDo(Ndara())
      | 'ShkruajNëBigQuery' >> beam.io.WriteToBigQuery('{0}:userlogs.logdata'.format(PROJEKTI), schema=schema,
        write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND)
   )
   rezultati = p.run()
   rezultati.wait_until_finish()

nëse __name__ == '__main__':
  logger = logging.getLogger().setLevel(logging.INFO)
  main()

Fillimi i konvejrit

Ne mund të nisim tubin në disa mënyra të ndryshme. Nëse do të donim, do të mund ta nisnim thjesht lokal nga terminali, duke hyrë larg në GCP.

python -m main_pipeline_stream.py 
 --input_topic "projekti/user-logs-237110/topikët/userlogs" 
 --streaming

Megjithatë, ne do ta nisim atë duke përdorur DataFlow. Ne mund ta bëjmë këtë me komandën e mëposhtme, duke vendosur parametrat e detyrueshëm të mëposhtëm.

  • projekt — ID e projektit tuaj GCP.
  • runner — mjeti qĂ« nxit tubin, i cili do tĂ« analizojĂ« programin tuaj dhe do tĂ« ndĂ«rtojĂ« tubin tuaj. PĂ«r tĂ« funksionuar nĂ« re, duhet tĂ« specifikoni DataflowRunner.
  • staging_location — rruga drejt ruajtjes nĂ« re Cloud Dataflow pĂ«r indeksimin e paketave tĂ« kodit qĂ« nevojiten pĂ«r pĂ«rpunuesit qĂ« kryejnĂ« punĂ«n.
  • temp_location — rruga drejt ruajtjes nĂ« re Cloud Dataflow pĂ«r ruajtjen e skedarĂ«ve pĂ«rkohĂ«sorĂ« tĂ« detyrave qĂ« krijohen gjatĂ« funksionimit tĂ« tubit.
  • streaming

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

Ndërsa ky komandë po ekzekutohet, mund të kalojmë në skedën DataFlow në google-console dhe të shqyrtojmë pipeline tonë. Duke klikuar në pipeline, duhet të shohim diçka të ngjashme me figurën 4. Për qëllime debuguese, mund të jetë shumë e dobishme të kalojmë në log-et dhe pastaj në Stackdriver për të parë log-et me detaje. Kjo më ndihmoi të zgjidhja probleme me pipeline në disa raste.

Krijojmë një kanal për përpunimin e të dhënave. Pjesa 2
Figura 4: Pipeline Beam

Qasja në të dhënat tona në BigQuery

Pra, tashmë duhet të kemi një pipeline të aktivizuar me të dhëna që po i japin tabelës tonë. Për ta kontrolluar këtë, mund të kalojmë te BigQuery dhe të shohim të dhënat. Pas përdorimit të komandës më poshtë, duhet të shihni disa rreshta të parë nga seti i të dhënave. Tani që kemi të dhëna të ruajtura në BigQuery, mund të bëjmë analiza të mëtejshme dhe gjithashtu të ndajmë të dhënat me kolegët dhe të fillojmë të përgjigjemi në pyetje biznesi.

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

Krijojmë një kanal për përpunimin e të dhënave. Pjesa 2
Figura 5: BigQuery

Përfundim

Shpresojmë që ky postim të shërbejë si një shembull i dobishëm për krijimin e një pipeline të të dhënave streaming dhe gjithashtu për gjetjen e mënyrave për ta bërë të dhënat më të aksesueshme. Ruajtja e të dhënave në një format të tillë na ofron shumë përfitime. Tani mund të fillojmë të përgjigjemi në pyetje të rëndësishme, siç janë: sa njerëz po e përdorin produktin tonë? A po rritet numri i përdoruesve me kalimin e kohës? Me cilat aspekte të produktit ndërveprojnë më shumë njerëzit? Dhe a ka gabime atje ku nuk duhet të jenë? Këto janë pyetje që do të jenë të rëndësishme për organizatën. Në bazë të ideve që rrjedhin nga përgjigjet për këto pyetje, do të jemi në gjendje të përmirësojmë produktin dhe të rrisim angazhimin e përdoruesve.

Beam është vërtet i dobishëm për këtë lloj ushtrimesh, si dhe ka një numër rastesh interesantë të tjera për përdorim. Për shembull, mund të analizoni të dhënat nga tick-at e bursës në kohë reale dhe të kryeni tregti në bazë të analizës, ndoshta keni të dhëna sensorësh që vijnë nga automjetet dhe dëshironi të llogaritni nivelin e trafikut. Po ashtu, mund të jeni një kompani lojërash që mbledh të dhëna nga përdoruesit dhe i përdor ato për të krijuar bordet informuese për të ndjekur indikatorët kyç. Mirë, zotërinj, kjo është një temë për një post tjetër, faleminderit për leximin dhe për ata që duan të shohin kodin e plotë, më poshtë është lidhja për GitHub-in tim.

https://github.com/DFoly/User_log_pipeline

Këtu përfundon gjithçka. Lexoni pjesën e parë.

Burimi: habr.com

Blini hosting tĂ« besueshĂ«m pĂ«r faqe interneti me mbrojtje nga DDoS, serverĂ« VPS VDS đŸ”„ Blini hosting tĂ« besueshĂ«m pĂ«r faqe interneti me mbrojtje nga DDoS, serverĂ« VPS VDS | ProHoster