Loome andmete voogude töötlemise toru. Osa 2

Tere kõigile. Jagame kursuse üliõpilastele spetsiaalselt valmistatud artikli viimase osa tõlget „Andmeinsener“. Esimese osaga saab tutvuda siin.

Apache Beam ja DataFlow reaalajas konveierite jaoks

Loome andmete voogude töötlemise toru. Osa 2

Google Cloudi seadistamine

Märkus: kasutades Google Cloud Shell'i andmete konveieri käivitamiseks ja kasutajalogide avaldamiseks, tekkis mul probleeme konveieri käivitamisega Python 3 peal. Google Cloud Shell kasutab Python 2, mis sobib paremini Apache Beami jaoks.

Konveieri käivitamiseks peame natuke seadistustes tuhnima. Teistest teistest GCP-del pole kasutanud, peate järgima järgmisi 6 sammu, mis on esitatud sellel lehelt.

Pärast seda peame meie skriptid Google'i pilvesalvestusse laadima ja neid oma Google Cloud Shell'i kopeerima. Laadimine pilvesalvestusse on piisavalt triviaalne (kirjeldus on saadaval siit). Meie failide kopeerimiseks saame avada Google Cloud Shell'i tööriistaribalt, klõpsates alloleval joonisel 2 vasakul esimesel ikoonil.

Loome andmete voogude töötlemise toru. Osa 2
Joonis 2

Failide kopeerimiseks ja vajalike raamatukogude installimiseks vajalikud käsklused on loetletud allpool.

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

Meie andmebaasi ja tabeli loomine

Kui oleme läbinud kõik seadistamise sammud, on järgmine asi, mida peame tegema, luua BigQuerys andmestik ja tabel. On mitu viisid, kuidas seda teha, kuid kõige lihtsam on kasutada Google Cloudi konsooli, alustades andmestiku loomisega. Saate järgida järgmisi lingi kaudu, et luua tabel oma skeemiga. Meie tabelil on 7 veergu, mis vastavad iga kasutajapõhise logi komponentidele. Mugavuse huvides määratleme kõik veerud stringidena (tüüp string), välja arvatud muutuja timelocal, ja nimetame need vastavalt varasemalt genereeritud muutujatele. Meie tabeli skeem peaks välja nägema nagu joonisel 3.

Loome andmete voogude töötlemise toru. Osa 2
Joonis 3. Tabeli skeem

Kasutajapõhise logi andmete avaldamine

Pub/Sub on meie voosüsteemi kriitiliselt tähtis komponent, mis võimaldab mitmel sõltumatul rakendusel omavahel suhelda. Eriti toimib see vahendajana, võimaldades meil rakenduste vahel sõnumeid saata ja vastu võtta. Esimene asi, mida peame tegema, on teema (topic) loomine. Piisab, kui minna Pub/Sub-i konsooli ja vajutada LOOMINE.

Allolev kood kutsub meie skripti, et genereerida logi andmeid, nagu varem määratletud, ning seejärel ühendub ja saadab logid Pub/Sub-i. Ainus, mida me peame tegema, on luua objekt PublisherClient, määrata teema rada meetodiga topic_path ja kutsuda välja funktsioon publish koos topic_path ja andmetega. Pange tähele, et impordime generate_log_line meie skriptist stream_logs, seega veenduge, et need failid on samas kaustas, vastasel juhul saate impordi tõrke. Siis saame seda käivitada meie google-konsoolis, kasutades:

python publish.py

import stream_logs
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):
    if message_future.exception(timeout=30):
        print('Publishing message on {} threw an Exception {}.'.format(
            topic_name, message_future.exception()))
    else:
        print(message_future.result())

if __name__ == '__main__':
    while True:
        line = stream_logs.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)

Kui fail on käivitatud, saame jälgida logiandmete väljundit konsoolis, nagu on näidatud alloleval joonisel. See skript töötab seni, kuni me ei kasuta CTRL+C, et see lõpetada.

Loome andmete voogude töötlemise toru. Osa 2
Joonis 4. Väljund publish_logs.py

Meie vooluahela koodi kirjutamine

Nüüd, kui kõik on valmis, saame alustada kõige huvitavamat osa — meie konveieri koodi kirjutamist, kasutades Beam'i ja Pythoni. Beam-konveieri loomiseks peame looma konveieri objekti (p). Pärast objekti loomist saame järjestikku rakendada mitmeid funktsioone, kasutades operaatorit pipe (|). Üldiselt näeb töövoog välja nagu alloleval joonisel.

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

Meie koodis loome kaks kohandatud funktsiooni. Funktsiooni regex_clean, mis skaneerib andmeid ja ekstraheerib vastava rea PATTERNS nimekirja põhjal, kasutades funktsiooni re.search. Funktsioon tagastab koma kaudu eraldatud rea. Kui te pole regulaarsete väljendite spetsialist, soovitan tutvuda selle õpikuga ja harjutada märkmikus, et koodi katsetada. Pärast seda määratleme kohandatud ParDo-funktsiooni nimega Split, mis on Beam-transformatsiooni variatsioon paralleelse töötlemise jaoks. Pythonis tehakse seda spetsiifilisel viisil — me peame looma klassi, mis pärib Beam’i DoFn klassist. Funktsioon Split võtab eelmisest funktsioonist parsetud stringi ja tagastab sõnastike loendi, mille võtmed vastavad meie BigQuery tabeli veergude nimedele. On mõned asjad, mida tuleb selle funktsiooni kohta märkida: pidin selle tööle saamiseks impordima datetime sisemiselt funktsiooni. Saime alguses faili impordimisel veateate, mis oli kummaline. See loend antakse seejärel funktsioonile WriteToBigQuery, mis lihtsalt lisab meie andmed tabelisse. Koodi Batch DataFlow Job'i ja Streaming DataFlow Job'i jaoks on esitatud allpool. Ainus erinevus partii ja voogesituse koodi vahel on see, et partii töötlemisel loeme CSV-d src_path, kasutades funktsiooni ReadFromText Beam'ist.

Batch DataFlow Job (partiis töötlemine)

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

PROJEKT='user-logs-237110'
skeem = '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):

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

   p.run()

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

Streaming DataFlow Job (voogesituse töötlemine)

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

Konveieri käivitamine

Saame käivitada konveieri mitmel erineval viisil. Kui soovime, saame selle lihtsalt käivitada kohalikult terminalist, sisenedes eemalt GCP-sse.

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

Kuid me kavatseme selle käivitada DataFlow abil. Saame seda teha alloleva käsu abil, seadistades järgmised kohustuslikud parameetrid.

  • project — teie GCP projekti ID.
  • runner — konveieri käivitamise vahend, mis analüüsib teie programmi ja konstruerib teie konveieri. Pilves töötamiseks peate määrama DataflowRunner.
  • staging_location — rada Cloud Dataflow pilve salvestusse, et indekseerida koodipakette, mida vajavad töötajad, kes tööd teevad.
  • temp_location — rada Cloud Dataflow pilve salvestusse ajutiste failide asukoha määramiseks, mis loodakse konveieri töö ajal.
  • streaming

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

Kuna see käsk on töös, saame minna google'i konsooli DataFlow vahekaardile ja vaadata meie toru. Torule klikkides peaksime nägema midagi sarnast joonisele 4. Tõrkeotsingu vajadusel võib olla väga kasulik minna logidesse ja seejärel Stackdriverisse detailsete logide vaatamiseks. See on aidanud mul mitmel korral toru probleeme lahendada.

Loome andmete voogude töötlemise toru. Osa 2
Joonis 4: Beam-toru

Juura meie andmetele BigQuery's

Seega peaks meil juba olema tööle pandud toru, kuhu andmed voolavad meie tabelisse. Selle kontrollimiseks saame minna BigQuery'sse ja vaadata andmeid. Pärast alltoodud käskluse kasutamist peaksite nägema andmekogumi esimesi ridasid. Nüüd, kui meil on andmed hoitud BigQuery's, saame teha edasist analüüsi, jagada andmeid kolleegidega ja hakata vastama äriküsimustele.

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

Loome andmete voogude töötlemise toru. Osa 2
Joonis 5: BigQuery

Kokkuvõte

Loodame, et see postitus annab kasuliku näite andmevoogude loomisest ning toob esile viise, kuidas andmeid kergemini kätte saada. Andmete selle vormingus hoidmine pakub meile mitmeid eeliseid. Nüüd saame hakata vastama olulistele küsimustele, nagu näiteks, kui palju inimesi meie toodet kasutab? Kas kasutajate arv kasvab ajas? Milliste toote aspektidega inimesed kõige rohkem suhtlevad? Ja kas on vigu seal, kus neid ei peaks olema? Need on küsimused, mis organisatsiooni jaoks olulised. Vastuste põhjal, mis tulenevad nendest küsimustest, saame toote täiustamiseks ja kasutajate huvi suurendamiseks häid ideid.

Beam on tõeliselt kasulik selliste harjutuste jaoks ning sellel on palju muid huvitavaid kasutusjuhtumeid. Näiteks võite analüüsida börsi tehingute andmeid reaalajas ja teha tehinguid analüüsi põhjal. Võib-olla on teil andurid, mis saadavad andmeid sõidukitest, ja soovite arvutada liiklusvoogu. Samuti võite olla mänguettevõte, kes kogub kasutajaandmeid ning kasutab neid ülevaatetablettide loomisel, et jälgida võtmemõõdikute täitmist. Noh, härrased, see on juba teine teema, tänan lugemise eest, ning neile, kes soovivad näha täielikku koodi, leiate allpool lingi minu GitHubi lehele.

https://github.com/DFoly/User_log_pipeline

Sellega on kõik. Loe esimest osa.

Allikas: habr.com

Osta usaldusväärne veebihosting DDoS kaitsega, VPS VDS serverid 🔥 Osta usaldusväärne veebihosting DDoS kaitsega, VPS VDS serverid | ProHoster